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

Получение данных

Получение (в терминах подсистемы — «загрузка») читает сообщения из темы Kafka и передаёт каждое из них процедуре прикладной обработки. Подсистема отвечает за распределение разделов по потокам, позиции чтения и завершение сеанса; что делать с сообщением — решает прикладной код.

Как проходит сеанс

Регламентное задание КафкаЗагрузка (или ручной запуск с формы шины)
  └─ КафкаСервер.Загрузить(Шина)
      ├─ определение списка читаемых разделов
      ├─ распределение разделов по потокам
      └─ в каждом потоке:
           ├─ ПередЗагрузкойРаздела            ← можно исключить раздел из сеанса
           ├─ позиционирование по сохранённым позициям
           ├─ цикл чтения: сообщение → процедура загрузки → фиксация позиций
           └─ ПослеЗагрузкиРаздела

Разделы и потоки

Если тема шины задана точным именем, читаются все её разделы. Если имя заканчивается на *, оно считается шаблоном: подсистема запрашивает у шлюза список подходящих тем и берёт разделы каждой из них. Когда не нашлось ни одной темы, сеанс завершается ошибкой — тема с опечаткой не должна выглядеть как «нет данных».

Разделы случайным, но равномерным образом распределяются между потоками; потоков не может быть больше, чем разделов. Каждый поток — отдельное фоновое задание со своим получателем в шлюзе. При единственном потоке чтение идёт в текущем сеансе.

Скорость приёма упирается в число разделов

Один раздел всегда читается одним потоком: это единственный способ сохранить порядок сообщений внутри раздела. Чтобы ускорить приём, увеличивают число разделов темы, а затем — число потоков шины.

Структура сообщения

В процедуру загрузки, указанную в реквизите шины «Процедура загрузки», передаётся один параметр — структура:

Поле Тип Значение
Тема Строка тема, из которой прочитано сообщение
Раздел Число номер раздела
Смещение Число смещение сообщения в разделе
ОтметкаВремени Число отметка времени сообщения в миллисекундах
Заголовки Соответствие заголовки сообщения; значение — строка или Null
Ключ зависит от сердеса ключ сообщения
Значение зависит от сердеса значение сообщения

Типы ключа и значения определяются сердесами шины. Пустое тело сообщения и сообщение-надгробие приходят по-разному в зависимости от формата — полная таблица приведена в разделе Сериализация ключей и значений.

// Модуль объекта обработки обмена
Процедура ПринятьСообщение(Сообщение) Экспорт

    Если Сообщение.Значение = Null Тогда
        // Надгробие: объект удалён на стороне отправителя.
        ПометитьНаУдаление(Сообщение.Ключ);
        Возврат;
    КонецЕсли;

    Если Сообщение.Значение = Неопределено Тогда
        // Пустое тело — для этого формата значения нет.
        Возврат; // или ВызватьИсключение
    КонецЕсли;

    ЗаписатьДанные(Сообщение.Ключ, Сообщение.Значение);

КонецПроцедуры

Транзакция не должна пережить сообщение

Процедура загрузки обязана завершать все начатые транзакции. Если после её вызова обнаружена активная транзакция, сеанс прерывается ошибкой «Обнаружена активная транзакция между получением сообщений»: иначе фиксация позиции чтения попала бы внутрь транзакции обработки данных и могла бы откатиться вместе с ней.

Позиции чтения

Позиции подсистема ведёт в информационной базе, а не в Kafka. Перед чтением раздела она берёт сохранённую позицию и переходит на неё; если раздел ещё не читался, чтение начинается с начала.

Сохранённая позиция сверяется с фактическими границами раздела и при необходимости корректируется — с предупреждением в журнале регистрации:

Ситуация Что произошло Как поправляется
позиция больше конечной тема пересоздана или очищена чтение продолжится с конца
позиция меньше начальной старые сообщения удалены по retention чтение продолжится с начала доступных

Позиции фиксируются по ходу чтения — периодически и по завершении обработки очередной порции сообщений, всегда после того, как сообщения приняты базой. Если процедура загрузки завершилась ошибкой, подсистема откатывает незавершённые транзакции и фиксирует позиции до сбойного сообщения: следующий сеанс начнёт ровно с него.

Гарантия приёма — «хотя бы один раз»

Сообщение считается принятым после успешного возврата из процедуры загрузки, но фиксация позиции происходит не на каждом сообщении. При обрыве связи или аварийном завершении сеанса часть уже принятых сообщений будет прочитана повторно. Приёмная логика должна быть идемпотентной.

Завершение сеанса

Обычный режим: сеанс завершается, когда все разделы дочитаны до конечных смещений, зафиксированных на старте, либо когда очередное обращение к шлюзу не вернуло сообщений. Данные, поступившие в тему уже во время сеанса, дождутся следующего запуска.

Режим «Приём в реальном времени»: сеанс не завершается на конце темы, а продолжает ждать новые сообщения и обрабатывать их сразу по мере поступления. Такой сеанс ограничен по времени — реквизит шины «Длительность сеанса» обязателен. Реальное время имеет смысл сочетать с частым расписанием: закончился один сеанс — регламентное задание тут же запускает следующий.

Независимо от режима, поток завершает работу досрочно, если:

  • истекла длительность сеанса;
  • завершился или был отменён родительский сеанс, запустивший потоки;
  • рабочий процесс сервера 1С требует завершения соединения (перезапуск, обновление).

Незавершённое чтение не теряется: позиции зафиксированы, и следующий сеанс продолжит с того же места.

События раздела

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