Загрузка данных


//++Варнин А.А. Задача №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");
	Исключение
		ЗаписьЖурналаРегистрации(ИмяСобытия, УровеньЖурналаРегистрации.Ошибка, , ,
			"Компонента подключилась, но объект не создался: "
				+ ПодробноеПредставлениеОшибки(ИнформацияОбОшибке()));
		Возврат Неопределено;
	КонецПопытки;

	Возврат Компонента;

КонецФункции

#КонецОбласти

#КонецОбласти