Позиции чтения¶
Позиция — это смещение следующего сообщения, которое обмен прочитает из раздела темы. Подсистема хранит позиции в информационной базе, в регистре сведений «Кафка: Позиции чтения», а не в Kafka.
| Поле регистра | Назначение |
|---|---|
| Период | момент фиксации, округлённый до часа |
| Кластер, Тема, Раздел | измерения: чей это раздел |
| Позиция | смещение следующего читаемого сообщения |
| Дата установки | точный момент фиксации |
Актуальной считается последняя запись по разделу. За час создаётся одна запись, по разделу хранится история последних записей — по ней видно, как двигалось чтение; лишние записи подсистема удаляет сама по завершении сеанса.
Почему позиции в базе, а не в Kafka¶
Штатный механизм Kafka фиксирует позиции автоматически, по таймеру, независимо от того, что произошло с данными на стороне получателя. Подсистема отключает автофиксацию (enable.auto.commit=false) и фиксирует позицию только после того, как сообщение принято информационной базой.
Отсюда два следствия:
- сообщения не теряются при обрыве связи или аварийном завершении: непринятое не будет отмечено как прочитанное;
- позиция не «уезжает» при перезапуске — она лежит в той же базе, что и принятые данные, и восстанавливается вместе с ней из резервной копии.
Именно поэтому параметры enable.auto.commit и auto.offset.reset не следует переопределять в настройках кластера или шины.
Где смотреть¶
Форма шины — основной инструмент: по каждому разделу видны зафиксированная позиция, дата фиксации, а по команде обновления — начальная и конечная позиции в Kafka и отставание. Отставание показывает, сколько сообщений ещё не прочитано.
Список регистра — полная история фиксаций по всем шинам и кластерам.
Рабочее место администратора показывает позиции групп получателей на стороне Kafka. Для обмена, который ведёт подсистема, они не используются — не путайте их с позициями регистра.
Перечитать или перемотать¶
Позиция — обычные данные, и её можно изменить: записать нужное смещение в регистр по кластеру, теме и разделу.
| Задача | Что сделать |
|---|---|
| перечитать тему с начала | удалить записи по разделам или установить позицию 0 |
| пропустить накопившееся | установить позицию, равную конечному смещению раздела |
| перечитать с определённого места | установить нужное смещение — его видно в просмотре сообщений темы |
Обмен должен быть остановлен
Правка позиции во время работающего сеанса бессмысленна: поток ведёт собственный счётчик и перезапишет позицию при следующей фиксации. Выключите обмен по расписанию, дождитесь завершения фоновых заданий и только потом меняйте позиции.
Повторное чтение — не ошибка, но и не бесплатно
Перечитывание темы с начала повторно вызовет процедуру загрузки для всех сообщений. Прикладная логика приёма должна быть идемпотентной — это требование и без того обязательно, поскольку гарантия приёма в подсистеме «хотя бы один раз».
Автоматическая коррекция¶
Перед чтением подсистема сверяет сохранённую позицию с фактическими границами раздела:
| Ситуация | Причина | Результат |
|---|---|---|
| позиция больше конечного смещения | тема пересоздана или очищена | чтение продолжится с конца, в журнал пишется предупреждение |
| позиция меньше начального смещения | старые сообщения удалены по retention | чтение продолжится с первого доступного, в журнал пишется предупреждение |
Такие предупреждения — не сбой обмена, но повод разобраться: регулярная коррекция «снизу» означает, что обмен не успевает за политикой хранения темы.