Перейти к содержанию

Прямая работа с Kafka

Обмен по расписанию покрывает типовой сценарий, но не ограничивает прикладной код: любой модуль может обратиться к Kafka напрямую — с теми же настройками подключения, что и настроенный обмен. Это удобно, когда нужно отправить сообщение по событию, прочитать ответ из служебной темы или выполнить административную операцию, не заводя отдельную шину.

Подключение

Кластер = КафкаСервер.КластерПоУмолчанию();          // или ссылка на нужный кластер
Адаптер = КафкаСервер.Адаптер(Кластер);              // HTTP-клиент шлюза

Адаптер — экземпляр обработки КафкаАдаптер с настроенным соединением. Он не обращается к Kafka сам по себе: сначала в шлюзе создаётся клиент.

Функция Что создаёт Базовая конфигурация
КафкаСервер.СоздатьProducer(Кластер, Адаптер, Операция, ДопКонфигурация, ТаймаутМс) отправителя конфигурация отправителя кластера
КафкаСервер.СоздатьConsumer(Кластер, Адаптер, Операция, ДопКонфигурация, ТаймаутМс) получателя конфигурация получателя кластера
КафкаСервер.СоздатьAdmin(Кластер, Адаптер, Операция, ДопКонфигурация, ТаймаутМс) администратора базовая конфигурация кластера
  • Операция — произвольное имя, под которым клиент виден в списке клиентов шлюза. Указывайте осмысленное: по нему разбирают, кто занял ресурсы.
  • ДопКонфигурация — соответствие с параметрами Kafka, перекрывающими конфигурацию кластера.
  • ТаймаутМс — время бездействия, после которого шлюз удалит клиента сам.

Обработка ошибок

Методы адаптера не выбрасывают исключений на ошибки шлюза и Kafka: они возвращают Неопределено, а описание остаётся в самом адаптере.

Producer = КафкаСервер.СоздатьProducer(Кластер, Адаптер, "Отправка уведомления");
Если Producer = Неопределено Тогда
    ВызватьИсключение Строка(Адаптер.КодОтвета) + " " + Адаптер.ОписаниеОшибки;
КонецЕсли;

Проверять результат нужно после каждого вызова — иначе ошибка проявится позже и в неочевидном месте.

Клиента нужно освободить

Клиент живёт в шлюзе и занимает соединение с Kafka. Освобождайте его в Попытка … Исключение, чтобы ресурс не остался висеть после ошибки. Таймаут бездействия — страховка, а не замена освобождению.

Отправка сообщения

Кластер = КафкаСервер.КластерПоУмолчанию();
Адаптер = КафкаСервер.Адаптер(Кластер);

Тема = КафкаСервер.Тема("events");

Producer = КафкаСервер.СоздатьProducer(Кластер, Адаптер, "Отправка уведомления");
Если Producer = Неопределено Тогда
    ВызватьИсключение Строка(Адаптер.КодОтвета) + " " + Адаптер.ОписаниеОшибки;
КонецЕсли;

Попытка

    Значение = КафкаСервер.JsonСериализовать(Данные);

    RecordMetadata = Адаптер.ProducerSend(Producer, Тема, Ключ, Значение);
    Если RecordMetadata = Неопределено Тогда
        ВызватьИсключение Строка(Адаптер.КодОтвета) + " " + Адаптер.ОписаниеОшибки;
    КонецЕсли;

    Адаптер.ProducerRelease(Producer);

Исключение
    Адаптер.ProducerRelease(Producer);
    ВызватьИсключение;
КонецПопытки;

ProducerSend передаёт ключ и значение строками. Функция КафкаСервер.Тема подставляет имя темы с учётом контура бэкапа — в копии рабочей базы сообщение уйдёт в тестовую тему.

Чтение сообщений

Тема = КафкаСервер.Тема("replies");

Consumer = КафкаСервер.СоздатьConsumer(Кластер, Адаптер, "Чтение ответов");
Если Consumer = Неопределено Тогда
    ВызватьИсключение Строка(Адаптер.КодОтвета) + " " + Адаптер.ОписаниеОшибки;
КонецЕсли;

Попытка

    Разделы = Адаптер.ConsumerGetPartitions(Consumer, Тема);
    // ... сформировать массив структур "topic, partition"
    Адаптер.ConsumerAssign(Consumer, Partitions);
    Адаптер.ConsumerSeekToEnd(Consumer, Partitions);

    Сообщения = Адаптер.ConsumerPoll(Consumer, 5000);

    Адаптер.ConsumerRelease(Consumer);

Исключение
    Адаптер.ConsumerRelease(Consumer);
    ВызватьИсключение;
КонецПопытки;

Получатель можно либо подписать на темы (ConsumerSubscribe — с распределением разделов внутри группы), либо назначить ему конкретные разделы (ConsumerAssign). Позиция задаётся методами ConsumerSeek, ConsumerSeekToBeginning, ConsumerSeekToEnd.

Позиции при прямом чтении

Регистр позиций чтения ведёт только автоматический обмен. При прямом чтении позицией управляет прикладной код — либо через ConsumerCommit (позиции хранит Kafka), либо самостоятельно.

Транзакции отправителя

ProducerBeginTransaction, ProducerCommitTransaction, ProducerAbortTransaction дают транзакционную отправку Kafka. Для этого при создании отправителя нужно передать в ДопКонфигурация параметр transactional.id. Транзакция Kafka не связана с транзакцией 1С — они фиксируются независимо.

Состав низкоуровневого API

Обработка КафкаАдаптер — тонкая обёртка над эндпойнтами шлюза: по методу на операцию, параметры и результат передаются структурами в терминах Kafka.

Группа Что доступно
Producer* создание и освобождение отправителя, отправка сообщений, разделы темы, транзакции
Consumer* подписка и назначение разделов, перемотка, чтение, фиксация позиций, метаданные группы, список тем
Admin* описание кластера и брокеров, темы и их разделы, конфигурации, группы получателей и их позиции, пользователи SCRAM, правила ACL, удаление записей
GetVersion версия шлюза

Полный перечень методов — в модуле объекта обработки КафкаАдаптер; имена совпадают с операциями шлюза, поэтому его документация применима напрямую.

Прикладной код обычно не создаёт адаптер вручную

В обработчике обмена готовый адаптер уже передаётся в события ПередВыгрузкойШины, ПередВыгрузкойПакета и ПередЗагрузкойРаздела — там достаточно воспользоваться им, не создавая собственный.