Загрузка данных
//++Варнин А.А. Задача №72300
// Общий модуль: кор_HTTPКоннектор
//
// Функция работает на любой из двух сборок компоненты, которые могут оказаться
// загруженными в рабочем процессе:
// 1.4.4 (сборка «Запчастей») — нет ОтправитьПакетСРезультатом,
// нет УстановитьТаймаутОчисткиПродюсера, ОстановитьПродюсера — процедура
// 1.9.1 (наш макет) — всё перечисленное есть, ОстановитьПродюсера — функция
//
// Набор возможностей определяется один раз при старте и пишется в журнал регистрации,
// чтобы после каждого прогона было видно, на какой сборке он отработал.
//
// ВАЖНО про проверку результата ОтправитьСообщение:
// 0 - сообщение поставлено в очередь (успех, асинхронный вариант)
// -1 - ошибка постановки в очередь
// Истина - отправлено (булев вариант на старых сборках)
// Поэтому проверка сделана через явный список успешных значений. Проверка вида
// «Если РезультатОтправки Тогда» на асинхронной сборке считает все нули ошибками,
// а проверка «= -1» на булевой сборке наоборот считает ошибки успехами.
#Область Кафка
#Область ПрограммныйИнтерфейс
Функция ПолучитьСтруктуруПодключения() Экспорт
СтруктураПодключения = Новый Структура;
СтруктураПодключения.Вставить("Брокеры", СокрЛП(Константы.кор_KafkaБрокеры.Получить()));
СтруктураПодключения.Вставить("ИмяТопика", СокрЛП(Константы.кор_KafkaТопик.Получить()));
СтруктураПодключения.Вставить("КаталогЛогов", СокрЛП(Константы.кор_KafkaКаталогЛогов.Получить()));
Возврат СтруктураПодключения;
КонецФункции
// Отправляет поток сообщений в Kafka.
//
// Параметры:
// Сообщения - Массив структур - каждый элемент: Ключ (Строка), Тело (Строка JSON)
//
// Возвращаемое значение:
// Структура:
// Отправлено - Число - принято компонентой без ошибок
// Ошибок - Число
// Успех - Булево
// ПодтверждениеДоставки - Булево - Ложь, если сборка не умеет подтверждать flush;
// тогда Отправлено означает «поставлено в очередь»
// Текст - Строка - готовая строка для журнала
//
Функция ОтправитьВКафку(Сообщения) Экспорт
ИмяСобытия = "Kafka. Отправка цен";
Результат = Новый Структура;
Результат.Вставить("Отправлено", 0);
Результат.Вставить("Ошибок", 0);
Результат.Вставить("Успех", Ложь);
Результат.Вставить("ПодтверждениеДоставки", Ложь);
Результат.Вставить("Текст", "");
Если ТипЗнч(Сообщения) <> Тип("Массив") ИЛИ Сообщения.Количество() = 0 Тогда
Результат.Текст = "Нет данных для отправки.";
Возврат Результат;
КонецЕсли;
Подключение = ПолучитьСтруктуруПодключения();
Если НЕ ЗначениеЗаполнено(Подключение.Брокеры) Тогда
Результат.Текст = "Не заполнена константа кор_KafkaБрокеры.";
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Возврат Результат;
КонецЕсли;
Если НЕ ЗначениеЗаполнено(Подключение.ИмяТопика) Тогда
Результат.Текст = "Не заполнена константа кор_KafkaТопик.";
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Возврат Результат;
КонецЕсли;
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Результат.Текст = "Не удалось создать компоненту Kafka.";
Возврат Результат;
КонецЕсли;
// --- Определяем, что умеет загруженная сборка -------------------------------
// ОстановитьПродюсера здесь не проверяем: вызов реально остановил бы продюсера.
Возможности = Новый Структура;
Возможности.Вставить("ПакетнаяОтправка", МетодДоступен(Компонента, "ОтправитьПакетСРезультатом"));
Возможности.Вставить("ТаймаутОчистки", МетодДоступен(Компонента, "УстановитьТаймаутОчисткиПродюсера"));
Возможности.Вставить("Метрики", МетодДоступен(Компонента, "ПолучитьМетрикиПродюсера"));
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Информация, , ,
СтрШаблон("Сборка компоненты: пакетная отправка %1, таймаут очистки %2, метрики %3.",
?(Возможности.ПакетнаяОтправка, "есть", "НЕТ"),
?(Возможности.ТаймаутОчистки, "есть", "НЕТ"),
?(Возможности.Метрики, "есть", "НЕТ")));
Результат.ПодтверждениеДоставки = Возможности.ПакетнаяОтправка;
// --- Настройка --------------------------------------------------------------
Если ЗначениеЗаполнено(Подключение.КаталогЛогов) Тогда
Попытка
Компонента.КаталогЛогов = Подключение.КаталогЛогов;
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Предупреждение, , ,
"Не удалось установить КаталогЛогов: "
+ ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
КонецПопытки;
КонецЕсли;
Если Возможности.ТаймаутОчистки Тогда
Компонента.УстановитьТаймаутОчисткиПродюсера(60000);
КонецЕсли;
Если НЕ Компонента.ИнициализироватьПродюсера(Подключение.Брокеры) Тогда
Результат.Текст = "Не удалось инициализировать продюсера. Брокеры: "
+ Подключение.Брокеры + ". " + Компонента.ПолучитьСообщениеОбОшибке();
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
Компонента = Неопределено;
Возврат Результат;
КонецЕсли;
Всего = Сообщения.Количество();
// --- Отправка ---------------------------------------------------------------
Попытка
ОтправитьПакетами(Компонента, Сообщения, Подключение.ИмяТопика,
Результат, Возможности, ИмяСобытия);
Исключение
Результат.Текст = "Прервано на " + (Результат.Отправлено + Результат.Ошибок)
+ " из " + Всего + ". " + ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , , Результат.Текст);
ОстановитьПродюсераБезопасно(Компонента, ИмяСобытия);
Компонента = Неопределено;
Возврат Результат;
КонецПопытки;
// --- Остановка. На 1.9.1 функция подтверждает flush, на 1.4.4 это процедура --
ОчередьОчищена = ОстановитьПродюсераБезопасно(Компонента, ИмяСобытия);
Если Возможности.Метрики Тогда
Попытка
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Информация, , ,
"Метрики продюсера: " + Компонента.ПолучитьМетрикиПродюсера());
Исключение
// Метрики — диагностика, из-за них падать незачем.
КонецПопытки;
КонецЕсли;
Компонента = Неопределено;
Результат.Успех = (Результат.Ошибок = 0) И ОчередьОчищена;
Результат.Текст = СтрШаблон("%1: %2, ошибок: %3, всего: %4.",
?(Результат.ПодтверждениеДоставки, "Доставлено", "Поставлено в очередь"),
Результат.Отправлено, Результат.Ошибок, Всего);
Если НЕ Результат.ПодтверждениеДоставки Тогда
Результат.Текст = Результат.Текст
+ " Сборка компоненты не подтверждает доставку — сверьте количество"
+ " принятых сообщений на стороне потребителя.";
КонецЕсли;
ЗаписьЖурналаРегистрации(ИмяСобытия,
?(Результат.Успех, УровеньЖурналаРегистрации.Информация,
УровеньЖурналаРегистрации.Предупреждение),
, , Результат.Текст);
Возврат Результат;
КонецФункции
#КонецОбласти
#Область СлужебныеПроцедурыИФункции
// Проверяет, есть ли у загруженной сборки указанный метод.
// Вызывает метод без параметров: если метода нет — платформа скажет
// «Метод объекта не обнаружен», если есть — ругнётся на параметры или отработает.
//
// НЕ применять к методам, у которых вызов без параметров что-то меняет.
//
Функция МетодДоступен(Компонента, ИмяМетода)
Попытка
Выполнить("Компонента." + ИмяМетода + "()");
Возврат Истина;
Исключение
Возврат Найти(КраткоеПредставлениеОшибки(ИнформацияОбОшибке()), "не обнаружен") = 0;
КонецПопытки;
КонецФункции
// Останавливает продюсера независимо от того, функция это или процедура.
// Возвращает Ложь только если остановка реально не удалась.
//
Функция ОстановитьПродюсераБезопасно(Компонента, ИмяСобытия)
Попытка
// 1.9.1: функция, Истина = вся очередь доставлена.
Возврат Компонента.ОстановитьПродюсера();
Исключение
// Провалились сюда либо потому, что это процедура, либо из-за реальной ошибки.
КонецПопытки;
Попытка
// 1.4.4: процедура. Подтверждения flush не даёт, считаем успехом.
Компонента.ОстановитьПродюсера();
Возврат Истина;
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Ошибка остановки продюсера: " + ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
Возврат Ложь;
КонецПопытки;
КонецФункции
Процедура ОтправитьПакетами(Компонента, Сообщения, Топик, Результат, Возможности, ИмяСобытия)
РазмерПачки = 1000;
Пачка = Новый Массив;
Для Каждого ЭлементСообщения Из Сообщения Цикл
Ключ = ?(ЭлементСообщения.Свойство("Ключ"), Строка(ЭлементСообщения.Ключ), "");
Тело = ?(ЭлементСообщения.Свойство("Тело"), Строка(ЭлементСообщения.Тело), "");
Если НЕ ЗначениеЗаполнено(Тело) Тогда
Результат.Ошибок = Результат.Ошибок + 1;
Продолжить;
КонецЕсли;
СообщениеКомпоненты = Новый Структура;
СообщениеКомпоненты.Вставить("message", Тело);
Если ЗначениеЗаполнено(Ключ) Тогда
СообщениеКомпоненты.Вставить("key", Ключ);
КонецЕсли;
Пачка.Добавить(СообщениеКомпоненты);
Если Пачка.Количество() >= РазмерПачки Тогда
ОтправитьОднуПачку(Компонента, Пачка, Топик, Результат, Возможности, ИмяСобытия);
Пачка = Новый Массив;
КонецЕсли;
КонецЦикла;
Если Пачка.Количество() > 0 Тогда
ОтправитьОднуПачку(Компонента, Пачка, Топик, Результат, Возможности, ИмяСобытия);
КонецЕсли;
КонецПроцедуры
Процедура ОтправитьОднуПачку(Компонента, Пачка, Топик, Результат, Возможности, ИмяСобытия)
// Ветвимся по заранее определённой возможности, а не по перехвату исключения:
// иначе переход на поштучную отправку происходит молча.
Если НЕ Возможности.ПакетнаяОтправка Тогда
ОтправитьПоштучно(Компонента, Пачка, Топик, Результат, ИмяСобытия);
Возврат;
КонецЕсли;
Запись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 примеров хватит, остальное в логе компоненты.
Если СтрокиОшибок.Количество() >= 50 Тогда
Прервать;
КонецЕсли;
КонецЦикла;
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Недоставленные сообщения:" + Символы.ПС + СтрСоединить(СтрокиОшибок, Символы.ПС));
КонецПроцедуры
// Поштучная отправка. Единственный доступный путь на сборке 1.4.4.
//
Процедура ОтправитьПоштучно(Компонента, Пачка, Топик, Результат, ИмяСобытия)
Для Каждого СообщениеКомпоненты Из Пачка Цикл
Ключ = "";
СообщениеКомпоненты.Свойство("key", Ключ);
Тело = СообщениеКомпоненты["message"];
РезультатОтправки = Компонента.ОтправитьСообщение(Тело, Топик, -1, Ключ);
// 0 — поставлено в очередь (асинхронная сборка), Истина — отправлено (булева).
Успешно = (РезультатОтправки = 0) ИЛИ (РезультатОтправки = Истина);
Если Успешно Тогда
Результат.Отправлено = Результат.Отправлено + 1;
Иначе
Результат.Ошибок = Результат.Ошибок + 1;
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Не отправлено. Ключ: " + Ключ + ". Результат: " + Строка(РезультатОтправки)
+ ". " + Компонента.ПолучитьСообщениеОбОшибке());
КонецЕсли;
КонецЦикла;
КонецПроцедуры
Функция СоздатьКомпонентуКафкаКлиент() Экспорт
ИмяСобытия = "Kafka. Подключение компоненты";
Попытка
Подключено = ПодключитьВнешнююКомпоненту(
"ОбщийМакет.SimpleKafkaAdapter64",
"кор_KafkaClient",
ТипВнешнейКомпоненты.Native,
ТипПодключенияВнешнейКомпоненты.Изолированно);
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
Возврат Неопределено;
КонецПопытки;
Если НЕ Подключено Тогда
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Не удалось подключить компоненту из общего макета SimpleKafkaAdapter64.");
Возврат Неопределено;
КонецЕсли;
Попытка
Компонента = Новый("AddIn.кор_KafkaClient.simpleKafka1C");
Исключение
ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
"Компонента подключилась, но объект не создался: "
+ ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
Возврат Неопределено;
КонецПопытки;
Возврат Компонента;
КонецФункции
#КонецОбласти
#КонецОбласти