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


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

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

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

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

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