Загрузка данных
// Отправляет поток сообщений в Kafka.
//
// Параметры:
// Сообщения - Массив структур - каждый элемент: Ключ (Строка), Тело (Строка JSON)
//
// Возвращаемое значение:
// Структура - Отправлено (Число), Ошибок (Число), Успех (Булево), Текст (Строка)
//
Функция ОтправитьВKafka(Сообщения) Экспорт
ИмяСобытия = "Kafka. Отправка цен";
Результат = Новый Структура("Отправлено, Ошибок, Успех, Текст", 0, 0, Ложь, "");
Если ТипЗнч(Сообщения) <> Тип("Массив") ИЛИ Сообщения.Количество() = 0 Тогда
Результат.Текст = "Нет данных для отправки.";
Возврат Результат;
КонецЕсли;
Подключение = ПолучитьСтруктуруПодключения();
Если НЕ ЗначениеЗаполнено(Подключение.Брокеры) Тогда
Результат.Текст = "Не заполнена константа кор_KafkaБрокеры.";
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Возврат Результат;
КонецЕсли;
Если НЕ ЗначениеЗаполнено(Подключение.ИмяТопика) Тогда
Результат.Текст = "Не заполнена константа кор_KafkaТопик.";
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Возврат Результат;
КонецЕсли;
Компонента = СоздатьКомпонентуКафкаКлиент();
Компонента.КаталогЛогов = Подключение.КаталогЛогов;
Если Компонента = Неопределено Тогда
Результат.Текст = "Не удалось создать компоненту Kafka.";
Возврат Результат;
КонецЕсли;
Попытка
Компонента.УстановитьТаймаутОчисткиПродюсера(60000);
Исключение
КонецПопытки;
Если НЕ Компонента.ИнициализироватьПродюсера(Подключение.Брокеры) Тогда
Результат.Текст = "Не удалось инициализировать продюсера. Брокеры: " + Подключение.Брокеры
+ ". " + Компонента.ПолучитьСообщениеОбОшибке();
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Компонента = Неопределено;
Возврат Результат;
КонецЕсли;
Всего = Сообщения.Количество();
Попытка
ОтправитьПакетами(Компонента, Сообщения, Подключение.ИмяТопика, Результат, ИмяСобытия);
Исключение
Результат.Текст = "Прервано на " + (Результат.Отправлено + Результат.Ошибок) + " из " + Всего
+ ". " + ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Попытка Компонента.ОстановитьПродюсера(); Исключение КонецПопытки;
Компонента = Неопределено;
Возврат Результат;
КонецПопытки;
ОчередьОчищена = ОстановитьПродюсераБезопасно(Компонента, ИмяСобытия);
Компонента = Неопределено;
Результат.Успех = (Результат.Ошибок = 0) И ОчередьОчищена;
Результат.Текст = СтрШаблон("Отправлено: %1, ошибок: %2, всего: %3.",
Результат.Отправлено, Результат.Ошибок, Всего);
ЗаписьЖурналаРегистрации(ИмяСобытия,
?(Результат.Успех, УровеньЖурналаРегистрации.Информация, УровеньЖурналаРегистрации.Предупреждение),
, , Результат.Текст);
Возврат Результат;
КонецФункции
Функция ОстановитьПродюсераБезопасно(Компонента, ИмяСобытия)
Попытка
Возврат Компонента.ОстановитьПродюсера();
Исключение
КонецПопытки;
Попытка
Компонента.ОстановитьПродюсера();
Возврат Истина;
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Ошибка остановки продюсера: " + ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
Возврат Ложь;
КонецПопытки;
КонецФункции
#КонецОбласти
#Область СлужебныеПроцедурыИФункции
Процедура ОтправитьПакетами(Компонента, Сообщения, Топик, Результат, ИмяСобытия)
РазмерПачки = 1000;
Пачка = Новый Массив;
Для Каждого ЭлементСообщения Из Сообщения Цикл
Ключ = ?(ЭлементСообщения.Свойство("Ключ"), Строка(ЭлементСообщения.Ключ), "");
Тело = ?(ЭлементСообщения.Свойство("Тело"), Строка(ЭлементСообщения.Тело), "");
Если НЕ ЗначениеЗаполнено(Тело) Тогда
Результат.Ошибок = Результат.Ошибок + 1;
Продолжить;
КонецЕсли;
СообщениеКомпоненты = Новый Структура;
СообщениеКомпоненты.Вставить("message", Тело);
Если ЗначениеЗаполнено(Ключ) Тогда
СообщениеКомпоненты.Вставить("key", Ключ);
КонецЕсли;
Пачка.Добавить(СообщениеКомпоненты);
Если Пачка.Количество() >= РазмерПачки Тогда
ОтправитьОднуПачку(Компонента, Пачка, Топик, Результат, ИмяСобытия);
Пачка = Новый Массив;
КонецЕсли;
КонецЦикла;
Если Пачка.Количество() > 0 Тогда
ОтправитьОднуПачку(Компонента, Пачка, Топик, Результат, ИмяСобытия);
КонецЕсли;
КонецПроцедуры
Процедура ОтправитьОднуПачку(Компонента, Пачка, Топик, Результат, ИмяСобытия)
ЗаписьJSON = Новый ЗаписьJSON;
ЗаписьJSON.УстановитьСтроку();
ЗаписатьJSON(ЗаписьJSON, Пачка);
ПачкаJSON = ЗаписьJSON.Закрыть();
ОтветJSON = "";
ЕстьПакетныйМетод = Истина;
Попытка
ОтветJSON = Компонента.ОтправитьПакетСРезультатом(ПачкаJSON, Топик, 60000);
Исключение
ЕстьПакетныйМетод = Ложь;
КонецПопытки;
Если НЕ ЕстьПакетныйМетод Тогда
ОтправитьПоштучно(Компонента, Пачка, Топик, Результат, ИмяСобытия);
Возврат;
КонецЕсли;
Если ПустаяСтрока(ОтветJSON) Тогда
Результат.Ошибок = Результат.Ошибок + Пачка.Количество();
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Пакет не отправлен: " + Компонента.ПолучитьСообщениеОбОшибке());
Возврат;
КонецЕсли;
ЧтениеJSON = Новый ЧтениеJSON;
ЧтениеJSON.УстановитьСтроку(ОтветJSON);
ОтветПакета = ПрочитатьJSON(ЧтениеJSON);
ЧтениеJSON.Закрыть();
Доставлено = Число(ОтветПакета["success_count"]);
НеДоставлено = Число(ОтветПакета["failed_count"]);
Результат.Отправлено = Результат.Отправлено + Доставлено;
Результат.Ошибок = Результат.Ошибок + НеДоставлено;
Если НеДоставлено > 0 Тогда
ЗаписатьОшибкиПакета(ОтветПакета, ИмяСобытия);
КонецЕсли;
КонецПроцедуры
Процедура ЗаписатьОшибкиПакета(ОтветПакета, ИмяСобытия)
СтрокиОшибок = Новый Массив;
Для Каждого СтрокаРезультата Из ОтветПакета["results"] Цикл
Если СтрокаРезультата["delivered"] = "true" Тогда
Продолжить;
КонецЕсли;
СтрокиОшибок.Добавить(СтрШаблон("Ключ: %1. %2",
СтрокаРезультата["key"], СтрокаРезультата["error"]));
Если СтрокиОшибок.Количество() >= 50 Тогда
Прервать;
КонецЕсли;
КонецЦикла;
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Недоставленные сообщения:" + Символы.ПС + СтрСоединить(СтрокиОшибок, Символы.ПС));
КонецПроцедуры
Процедура ОтправитьПоштучно(Компонента, Пачка, Топик, Результат, ИмяСобытия)
Для Каждого СообщениеКомпоненты Из Пачка Цикл
Ключ = "";
СообщениеКомпоненты.Свойство("key", Ключ);
Тело = СообщениеКомпоненты["message"];
РезультатОтправки = Компонента.ОтправитьСообщение(Тело, Топик, -1, Ключ);
Если РезультатОтправки = -1 Тогда
Результат.Ошибок = Результат.Ошибок + 1;
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Не поставлено в очередь. Ключ: " + Ключ + ". " + Компонента.ПолучитьСообщениеОбОшибке());
Иначе
Результат.Отправлено = Результат.Отправлено + 1;
КонецЕсли;
КонецЦикла;
КонецПроцедуры
Функция СоздатьКомпонентуКафкаКлиент() Экспорт
ИмяСобытия = "Kafka. Подключение компоненты";
Попытка
Подключено = ПодключитьВнешнююКомпоненту(
"ОбщийМакет.SimpleKafkaAdapter",
"кор_KafkaClient",
ТипВнешнейКомпоненты.Native,
ТипПодключенияВнешнейКомпоненты.Изолированно);
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
Возврат Неопределено;
КонецПопытки;
Если НЕ Подключено Тогда
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Не удалось подключить компоненту из общего макета SimpleKafkaAdapter64.");
Возврат Неопределено;
КонецЕсли;
Попытка
Компонента = Новый("AddIn.кор_KafkaClient.simpleKafka1C");
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Компонента подключилась, но объект не создался: "
+ ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
Возврат Неопределено;
КонецПопытки;
Возврат Компонента;
КонецФункции
#КонецОбласти
#КонецОбласти