Загрузка данных
// https://github.com/NuclearAPK/Simple-Kafka_Adapter
#Область Кафка
// Функция - Создать компоненту кафка клиент
//
// Возвращаемое значение:
// ТипВнешнейКомпоненты.Native, Неопределено - Внешняя компонента
//
Функция СоздатьКомпонентуКафкаКлиент() Экспорт
Если ПодключитьВнешнююКомпоненту("ОбщийМакет.SimpleKafkaAdapter64", "KafkaClient", ТипВнешнейКомпоненты.Native) Тогда
Попытка
Возврат Новый("AddIn.KafkaClient.simpleKafka1C");
Исключение
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Компонента из общего макета SimpleKafkaAdapter64 подключилась, но не создалась",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат Неопределено;
КонецПопытки;
Иначе
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Не удалось подключить компоненту из общего макета SimpleKafkaAdapter64",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат Неопределено;
КонецЕсли;
КонецФункции
// Процедура - Отправить в кафку
//
// Параметры:
// СтруктураПодключения - Структура - полученная через
// РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения("Наименование настройки")
// Данные - Строка - строка в формате Json
//
Процедура ОтправитьВКафку(СтруктураПодключения, Данные) Экспорт
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
Если ТипЗнч(Данные) <> Тип("Массив") Тогда
Сообщения = ОбщегоНазначенияКлиентСервер.ЗначениеВМассиве(Данные);
Иначе
Сообщения = Данные;
КонецЕсли;
//Если ТипЗнч(Параметры) = Тип("Соответствие") ИЛИ (ТипЗнч(Параметры) = Тип("Структура")) Тогда
// Для Каждого Параметр Из Параметры Цикл
// Компонента.УстановитьПараметр(Строка(Параметр.Ключ), Строка(Параметр.Значение));
// КонецЦикла;
//КонецЕсли;
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокеры = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топик = СтруктураПодключения.ИмяТопика;
// инициализируем подключение к брокеру
// достаточно указать одного брокера из кластера, либо есть возможность перечислить брокеров через ,
РезультатИнициализации = Компонента.ИнициализироватьПродюсера(Брокеры);
//РезультатИнициализации = Компонента.ИнициализироватьПродюсера("89.248.195.244:50086,89.248.195.244:50087,89.248.195.244:50088");
Если РезультатИнициализации Тогда
Для Каждого СообщениеВКафку Из Сообщения Цикл
// Parametr1 - Тело сообщения (строка)
// Parametr2 - Топик (строка)
// Parametr3 - Номер партиции, по умолчанию = -1 (число)
// Parametr4 - Произвольный ключ, идентифицирующий сообщение, например GUID (строка)
// Parametr5 - Заголовки (строка), "ключ1,значение1;ключ2,значение2"
//ИДТранзакции = Строка(Новый УникальныйИдентификатор);
//Заголовки = ""; //"key1,test1;key2,test2";
РезультатОтправки = Компонента.ОтправитьСообщение(СообщениеВКафку, Топик, 0);
Если Не РезультатОтправки Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(Компонента.ПолучитьСообщениеОбОшибке(),
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
КонецЦикла;
Компонента.ОстановитьПродюсера();
Компонента = Неопределено;
КонецЕсли;
КонецПроцедуры
// Процедуру запускает регламентное задание ЗапуститьЧитателяКафкиСсылкиНаОплату
Процедура ЗапуститьЧитателяКафкиСсылкиНаОплату() Экспорт
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеСсылокНаОплату";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
Компонента.УстановитьТаймаутОжидания(5000); // установка таймаута для ожидания сообщений - 5 сек.
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Если НЕ СтруктураПодключения.Параметры.РазрешеноСлушать Тогда
Прервать;
КонецЕсли;
Попытка
ЕстьСообщение = Компонента.ПрочитатьСообщение();
Если Не ЕстьСообщение Тогда
Продолжить;
КонецЕсли;
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Сообщение = Компонента.ПолучитьДанныеСообщения();
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураДанных = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Топик = Компонента.ПолучитьТопикСообщения();
Смещение = Компонента.ПолучитьСмещениеСообщения();
Партиция = Компонента.ПолучитьРазделСообщения();
Если СтруктураДанных.Свойство("order_num") И ЗначениеЗаполнено(СтруктураДанных.order_num) Тогда
ЗаписатьСсылкуНаОплатуВЗаказКлиента(СтруктураДанных);
Иначе
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Пустое значение order_num",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(Топик, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
Процедура ЗаписатьСсылкуНаОплатуВЗаказКлиента(СтруктураДанных)
ДокСсылка = Документы.ЗаказДляКлиента.НайтиПоНомеру(СтруктураДанных.order_num);
Если Не ДокСсылка.Пустая() Тогда
ДокОбъект = ДокСсылка.ПолучитьОбъект();
НоваяСтрокаПлатежи = ДокОбъект.Платежи.Добавить();
НоваяСтрокаПлатежи.Дата = ТекущаяДатаСеанса();
Если СтруктураДанных.response.status = "error" Тогда
НоваяСтрокаПлатежи.СсылкаНаПлатеж = СтруктураДанных.response.message;
Иначе
НоваяСтрокаПлатежи.СсылкаНаПлатеж = СтруктураДанных.response.link;
КонецЕсли;
НоваяСтрокаПлатежи.Сумма = СтруктураДанных.sum;
Если СтруктураДанных.type = "partial" Тогда
Признак = "Предоплата";
ИначеЕсли СтруктураДанных.type = "full" Тогда
Признак = "ПолнаяПредоплата";
ИначеЕсли СтруктураДанных.type = "surcharge" Тогда
Признак = "Доплата";
Иначе
Признак = "Не определен";
КонецЕсли;
НоваяСтрокаПлатежи.Признак = Признак;
ДокОбъект.Записать();
Иначе
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Не найден документ с номером " + СтруктураДанных.order_num,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
КонецПроцедуры
Процедура ЗапуститьЧитателяКафкиПланыДашборд() Экспорт
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеПлановДашборд";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
//Топики = СтруктураПодключения.Параметры.ИмяТопика;
Топики = СтруктураПодключения.ИмяТопика;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Если НЕ СтруктураПодключения.Параметры.РазрешеноСлушать Тогда
Прервать;
КонецЕсли;
Попытка
Сообщение = Компонента.Слушать();
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
//Компонента = Неопределено;
//Возврат;
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураJSON = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
// ВременнаяМетка = Число(СтруктураJSON.timestamp);
Смещение = Число(СтруктураJSON.offset);
Партиция = Число(СтруктураJSON.partition);
Чтение = Новый ЧтениеJSON;
Чтение.УстановитьСтроку(СтруктураJSON.message);
СтруктураДанных = ПрочитатьJSON(Чтение);
кор_ДашбордыОбменКафка.СоздатьДокументПланаИзСтруктуры(СтруктураДанных);
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
Процедура ЗапуститьЧитателяКафкиЗН() Экспорт
//++НП / 29.03.2024 / № 45716 Дашборд
Если НЕ Константы.НоваяМотивацияВключена.Получить() Тогда
Возврат;
КонецЕсли;
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеЗНДашборд";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
//СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения_2("_ЧтениеКафка");
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
//Компонента.УстановитьПозициюЧтения("abcp-to-1c-dev", 5,0);
//Если корПолучитьЗначениеРеквизита("_ЧитаемЗаказыСКафки", ПланыВидовХарактеристик.ДопРеквизиты.Настройка) = Истина Тогда
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
//КонецЕсли;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Если НЕ СтруктураПодключения.Параметры.РазрешеноСлушать Тогда
Прервать;
КонецЕсли;
Попытка
ЕстьСообщение = Компонента.ПрочитатьСообщение();
Если Не ЕстьСообщение Тогда
Продолжить;
КонецЕсли;
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Сообщение = Компонента.ПолучитьДанныеСообщения();
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураДанных = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Топик = Компонента.ПолучитьТопикСообщения();
Смещение = Компонента.ПолучитьСмещениеСообщения();
Партиция = Компонента.ПолучитьРазделСообщения();
кор_ДашбордыОбменКафка.СоздатьЗаказ(СтруктураДанных);
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
//Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
РеквизитыОбъекта = РегистрыСведений.ДопРеквизиты.СоздатьНаборЗаписей();
РеквизитыОбъекта.Отбор.ТипОбъекта.Установить(ИмяСервиса);
РеквизитыОбъекта.Отбор.Реквизит.Установить(ПланыВидовХарактеристик.ДопРеквизиты.КафкаОффсет);
РеквизитыОбъекта.Прочитать();
Если РеквизитыОбъекта.Количество() > 0 Тогда
РеквизитыОбъекта[0].Значение = Смещение + 1;
РеквизитыОбъекта.Записать();
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(Топик, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
Процедура ЗапуститьЧитателяКафкиПолучитьОплатуДашборд() Экспорт
//++НП / 29.03.2024 / № 45716 Дашборд
Если НЕ Константы.НоваяМотивацияВключена.Получить() Тогда
Возврат;
КонецЕсли;
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеЗНПолучитьОплатуДашборд";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
//СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения_2("_ЧтениеКафка");
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
//Компонента.УстановитьПозициюЧтения("abcp-to-1c-dev", 5,0);
//Если корПолучитьЗначениеРеквизита("_ЧитаемЗаказыСКафки", ПланыВидовХарактеристик.ДопРеквизиты.Настройка) = Истина Тогда
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
//КонецЕсли;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
Попытка
ЕстьСообщение = Компонента.ПрочитатьСообщение();
Если Не ЕстьСообщение Тогда
Продолжить;
КонецЕсли;
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Сообщение = Компонента.ПолучитьДанныеСообщения();
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураДанных = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Топик = Компонента.ПолучитьТопикСообщения();
Смещение = Компонента.ПолучитьСмещениеСообщения();
Партиция = Компонента.ПолучитьРазделСообщения();
кор_ДашбордыОбменКафка.ПолучитьОплатуПоЗН(СтруктураДанных);
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
//Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
РеквизитыОбъекта = РегистрыСведений.ДопРеквизиты.СоздатьНаборЗаписей();
РеквизитыОбъекта.Отбор.ТипОбъекта.Установить(ИмяСервиса);
РеквизитыОбъекта.Отбор.Реквизит.Установить(ПланыВидовХарактеристик.ДопРеквизиты.КафкаОффсет);
РеквизитыОбъекта.Прочитать();
Если РеквизитыОбъекта.Количество() > 0 Тогда
РеквизитыОбъекта[0].Значение = Смещение + 1;
РеквизитыОбъекта.Записать();
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(Топик, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
Процедура ЗапуститьЧитателяКафкиДатаОтчетаПоПродажам() Экспорт
//++НП / Дашборд директора
Если НЕ Константы.НоваяМотивацияВключена.Получить() Тогда
Возврат;
КонецЕсли;
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеДатыДляОтчетаПоПродажам";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
//СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения_2("_ЧтениеКафка");
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
//Компонента.УстановитьПозициюЧтения("abcp-to-1c-dev", 5,0);
//Если корПолучитьЗначениеРеквизита("_ЧитаемЗаказыСКафки", ПланыВидовХарактеристик.ДопРеквизиты.Настройка) = Истина Тогда
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
//КонецЕсли;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
Попытка
ЕстьСообщение = Компонента.ПрочитатьСообщение();
Если Не ЕстьСообщение Тогда
Продолжить;
КонецЕсли;
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Сообщение = Компонента.ПолучитьДанныеСообщения();
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураДанных = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Топик = Компонента.ПолучитьТопикСообщения();
Смещение = Компонента.ПолучитьСмещениеСообщения();
Партиция = Компонента.ПолучитьРазделСообщения();
кор_ДашбордыОбменКафка.ПолучитьДатуДляФормированияОтчета(СтруктураДанных);
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
//Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
РеквизитыОбъекта = РегистрыСведений.ДопРеквизиты.СоздатьНаборЗаписей();
РеквизитыОбъекта.Отбор.ТипОбъекта.Установить(ИмяСервиса);
РеквизитыОбъекта.Отбор.Реквизит.Установить(ПланыВидовХарактеристик.ДопРеквизиты.КафкаОффсет);
РеквизитыОбъекта.Прочитать();
Если РеквизитыОбъекта.Количество() > 0 Тогда
РеквизитыОбъекта[0].Значение = Смещение + 1;
РеквизитыОбъекта.Записать();
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(Топик, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
//rialist+++
Процедура ЗапуститьЧитателяКафкиПолучитьСтатусОплатыЗК() Экспорт
//++НП / 29.03.2024 / № 45716 Дашборд
Если НЕ Константы.ВключенИнвестПроект2.Получить() Тогда
Возврат;
КонецЕсли;
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеЗНПолучитьСтатусОплаты";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
//СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения_2("_ЧтениеКафка");
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
//Компонента.УстановитьПозициюЧтения("abcp-to-1c-dev", 5,0);
//Если корПолучитьЗначениеРеквизита("_ЧитаемЗаказыСКафки", ПланыВидовХарактеристик.ДопРеквизиты.Настройка) = Истина Тогда
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
//КонецЕсли;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
Попытка
ЕстьСообщение = Компонента.ПрочитатьСообщение();
Если Не ЕстьСообщение Тогда
Продолжить;
КонецЕсли;
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Сообщение = Компонента.ПолучитьДанныеСообщения();
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураДанных = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Топик = Компонента.ПолучитьТопикСообщения();
Смещение = Компонента.ПолучитьСмещениеСообщения();
Партиция = Компонента.ПолучитьРазделСообщения();
кор_ДашбордыОбменКафка.ПолучитьСтатусОплатыПоЗК(СтруктураДанных);
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
//Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
РеквизитыОбъекта = РегистрыСведений.ДопРеквизиты.СоздатьНаборЗаписей();
РеквизитыОбъекта.Отбор.ТипОбъекта.Установить(ИмяСервиса);
РеквизитыОбъекта.Отбор.Реквизит.Установить(ПланыВидовХарактеристик.ДопРеквизиты.КафкаОффсет);
РеквизитыОбъекта.Прочитать();
Если РеквизитыОбъекта.Количество() > 0 Тогда
РеквизитыОбъекта[0].Значение = Смещение + 1;
РеквизитыОбъекта.Записать();
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(Топик, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
//---rialist
//Спиридонов АИ 45970
Процедура ЗапуститьЧитателяКафкиЗаказыABCP() Экспорт
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_ЧитаемЗаказыСКафки";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
//СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения_2("_ЧтениеКафка");
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
//Компонента.УстановитьПозициюЧтения("abcp-to-1c-dev", 5,0);
//Если корПолучитьЗначениеРеквизита("_ЧитаемЗаказыСКафки", ПланыВидовХарактеристик.ДопРеквизиты.Настройка) = Истина Тогда
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
//КонецЕсли;
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
Попытка
ЕстьСообщение = Компонента.ПрочитатьСообщение();
Если Не ЕстьСообщение Тогда
Продолжить;
КонецЕсли;
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Сообщение = Компонента.ПолучитьДанныеСообщения();
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураДанных = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Топик = Компонента.ПолучитьТопикСообщения();
Смещение = Компонента.ПолучитьСмещениеСообщения();
Партиция = Компонента.ПолучитьРазделСообщения();
Если СтруктураДанных.Свойство("number") И ЗначениеЗаполнено(СтруктураДанных.number) Тогда
Попытка
рсОбменССайтомABCP.ВыполнитьОбменЗаказамиИзКафки(СтруктураДанных, Смещение);
Исключение
рсОбменССайтомABCP.СообщитьОПроблемеЗаписиЗаказаСАВСР(СтруктураДанных.number, Строка(ОписаниеОшибки()));
КонецПопытки;
Иначе
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Пустое значение number",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
//Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
РеквизитыОбъекта = РегистрыСведений.ДопРеквизиты.СоздатьНаборЗаписей();
РеквизитыОбъекта.Отбор.ТипОбъекта.Установить(ИмяСервиса);
РеквизитыОбъекта.Отбор.Реквизит.Установить(ПланыВидовХарактеристик.ДопРеквизиты.КафкаОффсет);
РеквизитыОбъекта.Прочитать();
Если РеквизитыОбъекта.Количество() > 0 Тогда
РеквизитыОбъекта[0].Значение = Смещение + 1;
РеквизитыОбъекта.Записать();
КонецЕсли;
// вручную фиксируем смещения
// !!! второй параметр - смещение, должно быть на 1 больше текушего прочитанного смещения
Компонента.ЗафиксироватьСмещение(Топик, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
Процедура ЗапуститьЧитателяКафкиВерификацииДисконтныхКарт() Экспорт
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеВерификацииДисконтныхКарт";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
//Компонента.УстановитьПараметр("statistics.interval.ms", 15000);
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
// установка таймаута для ожидания сообщений - 5 сек.
Компонента.УстановитьТаймаутОжидания(5000);
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
Попытка
Сообщение = Компонента.Слушать();
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураJSON = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
// ВременнаяМетка = Число(СтруктураJSON.timestamp);
Смещение = Число(СтруктураJSON.offset);
Партиция = Число(СтруктураJSON.partition);
Чтение = Новый ЧтениеJSON;
Чтение.УстановитьСтроку(СтруктураJSON.message);
СтруктураДанных = ПрочитатьJSON(Чтение);
Если СтруктураДанных.Свойство("code_id") И ЗначениеЗаполнено(СтруктураДанных.code_id) Тогда
ЗаписатьВерификацииДисконтныхКарт(СтруктураДанных);
Иначе
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Пустое значение code_id",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
// вручную фиксируем смещения
РеквизитыОбъекта = РегистрыСведений.ДопРеквизиты.СоздатьНаборЗаписей();
РеквизитыОбъекта.Отбор.ТипОбъекта.Установить(ИмяСервиса);
РеквизитыОбъекта.Отбор.Реквизит.Установить(ПланыВидовХарактеристик.ДопРеквизиты.КафкаОффсет);
РеквизитыОбъекта.Прочитать();
Если РеквизитыОбъекта.Количество() > 0 Тогда
РеквизитыОбъекта[0].Значение = Смещение + 1;
РеквизитыОбъекта.Записать();
КонецЕсли;
// вручную фиксируем смещения
Компонента.ЗафиксироватьСмещение(СтруктураJSON.topic, Смещение + 1, Партиция);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
Процедура ЗаписатьВерификацииДисконтныхКарт(СтруктураДанных)
Интервал = 1800; //Срок Действия
Если СтруктураДанных.status = 0 Тогда
ВерификацияНабор = РегистрыСведений.ВерификацииДисконтныхКарт.СоздатьНаборЗаписей();
ВерификацияНабор.Отбор.Code_id.Установить(СтруктураДанных.code_id);
ВерификацияНабор.Прочитать();
Если ВерификацияНабор.Количество() > 0 Тогда
ВерификацияНабор[0].ПризнакВерификации = Истина;
КонецЕсли;
Попытка
ВерификацияНабор.Записать();
Исключение
ТекстОшибки = ОписаниеОшибки();
ЗаписьЖурналаРегистрации("В РС ВерификацииДисконтныхКарт запись неудалась ",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецПопытки;
ИначеЕсли СтруктураДанных.status = 1 Тогда
МенеджерЗаписи = РегистрыСведений.ВерификацииДисконтныхКарт.СоздатьМенеджерЗаписи();
Пробел = 160;
МенеджерЗаписи.ДатаСоздания = ИнтеграцияЕБК.ПолучитьДату(СтруктураДанных.created);
МенеджерЗаписи.ДатаДействия = ИнтеграцияЕБК.ПолучитьДату(СтруктураДанных.created) + Интервал;
МенеджерЗаписи.КодКлиента = СтрЗаменить(СтруктураДанных.ebk_id, Символ(Пробел), "");
МенеджерЗаписи.Код = СтрЗаменить(СтруктураДанных.code, Символ(Пробел), "");
МенеджерЗаписи.Code_id = СтруктураДанных.code_id;
МенеджерЗаписи.mobile_user_id = СтруктураДанных.mobile_user_id;
Попытка
МенеджерЗаписи.Записать();
Исключение
ТекстОшибки = ОписаниеОшибки();
ЗаписьЖурналаРегистрации("В РС ВерификацииДисконтныхКарт запись неудалась ",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецПопытки;
КонецЕсли;
КонецПроцедуры
Функция ПризнакВерифмкацииОтправитьВКафку(КодИД, ИмяСервиса) Экспорт
Статус = Ложь;
СтруктураОтправки = Новый Структура;
СтруктураОтправки.Вставить("Code_id", КодИД);
СтруктураОтправки.Вставить("status", 0);
Попытка
Джсон = КоннекторHTTP.ОбъектВJson(СтруктураОтправки);
Статус = Истина;
Исключение
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Верификация отправка Кафка",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецПопытки;
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
КоннекторHTTP.ОтправитьВКафку(СтруктураПодключения,СтрЗаменить(Джсон, Символы.ПС,""));
Возврат Статус;
КонецФункции
Процедура ЗапуститьЧитателяКафкиСделки() Экспорт
Компонента = СоздатьКомпонентуКафкаКлиент();
Если Компонента = Неопределено Тогда
Возврат;
КонецЕсли;
ИмяСервиса = "_КафкаЧтениеСделок";
СтруктураПодключения = РегламентныхЗаданийОбщий.ПолучитьСтруктуруПодключения(ИмяСервиса);
Компонента.УстановитьПараметр("group.id", СтруктураПодключения.Параметры.ГруппаПодписчика);
Компонента.УстановитьПараметр("enable.auto.commit", "false");
Компонента.УстановитьПараметр("enable.auto.offset.store", "false");
Компонента.УстановитьПараметр("enable.partition.eof", "false");
Если СтруктураПодключения.Параметры.Свойство("КаталогЛогов")
И ЗначениеЗаполнено(СтруктураПодключения.Параметры.КаталогЛогов) Тогда
Компонента.КаталогЛогов = СтруктураПодключения.Параметры.КаталогЛогов;
КонецЕсли;
Брокер = СтрЗаменить(СтруктураПодключения.Хост, "http://", "");
Топики = СтруктураПодключения.ИмяТопика;
Компонента.УстановитьПозициюЧтения(Топики, Цел(СтруктураПодключения.КафкаОффсет), 0);
Результат = Компонента.ИнициализироватьКонсьюмера(Брокер, Топики);
Компонента.УстановитьТаймаутОжидания(5000);
ЗаписьЖурналаРегистрации("УстановитьПозициюЧтения " + Строка(Топики),
УровеньЖурналаРегистрации.Информация, , ,
Топики + " " + Строка(СтруктураПодключения.КафкаОффсет));
Если Не Результат Тогда
ТекстОшибки = СтрШаблон("Не удалось инициализировать консьюмера для топиков: %1", Топики);
ЗаписьЖурналаРегистрации("Интеграция Kafka. Consumer",
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
Возврат;
КонецЕсли;
Пока Истина Цикл
Попытка
Сообщение = Компонента.Слушать();
Исключение
Компонента = Неопределено;
ВызватьИсключение ОписаниеОшибки();
КонецПопытки;
Если НЕ ЗначениеЗаполнено(Сообщение) Тогда
ОписаниеОшибкиКомпонента = Компонента.ПолучитьСообщениеОбОшибке();
Если ЗначениеЗаполнено(ОписаниеОшибкиКомпонента) Тогда
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации(ОписаниеОшибкиКомпонента,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецЕсли;
Продолжить;
КонецЕсли;
ОбъектЧтение = Новый ЧтениеJSON;
ОбъектЧтение.УстановитьСтроку(Сообщение);
СтруктураJSON = ПрочитатьJSON(ОбъектЧтение);
ОбъектЧтение.Закрыть();
Смещение = Число(СтруктураJSON.offset);
Партиция = Число(СтруктураJSON.partition);
Топик = СтруктураJSON.topic;
Чтение = Новый ЧтениеJSON;
Чтение.УстановитьСтроку(СтруктураJSON.message);
СтруктураДанных = ПрочитатьJSON(Чтение);
ЗаписьЖурналаРегистрации("ЧтениеСообщенияКафка " + Строка(Топик),
УровеньЖурналаРегистрации.Информация, , ,
"Смещение " + Строка(Смещение) +
" Партиция " + Строка(Партиция) +
" Топик " + Строка(Топик) +
" Сообщение " + Строка(Сообщение));
ПроверкиПройдены = ВоронкиПродажСервер.ПроверкаСтруктурыОбращения(Компонента, Топик, Смещение, Партиция, ИмяСервиса, СтруктураДанных);
Если Не ПроверкиПройдены Тогда
Продолжить;
КонецЕсли;
Попытка
ВоронкиПродажСервер.СоздатьДокументОбращениеВоронкиУдленногоОП(СтруктураДанных);
Исключение
ТекстОшибки = ПодробноеПредставлениеОшибки(ИнформацияОбОшибке());
ЗаписьЖурналаРегистрации("Ошибка создания обращения из сделки через кафку, оффсет: " + Смещение,
УровеньЖурналаРегистрации.Ошибка, , ,
ТекстОшибки);
КонецПопытки;
ВоронкиПродажСервер.ЗафиксироватьСмешение (Компонента, Топик, Смещение, Партиция, ИмяСервиса);
КонецЦикла;
Компонента.ОстановитьКонсьюмера();
Компонента = Неопределено;
КонецПроцедуры
#КонецОбласти