сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Kafka и стриминг

Гарантии доставки в Kafka: at-least-once и exactly-once

Семантики доставки

At-most-once, at-least-once, exactly-once - три обещания о том, что случится с сообщением при сбоях. Собес проверяет, понимаешь ли ты, что честный exactly-once дорог и упирается во внешний sink.

Стержень: дефолт индустрии - at-least-once плюс идемпотентный потребитель, а не мнимый exactly-once через границы систем.

// Формулировки: «какие семантики доставки бывают?», «бывает ли exactly-once на самом деле?», «где рождаются дубли?».

Три семантики и дефолт

At-most-once - можем потерять, но не задублируем. At-least-once - не потеряем, но возможны дубли. Exactly-once - не потеряем и не задублируем, но это дорого и с оговорками.

Дефолт индустрии - at-least-once плюс идемпотентный потребитель: дубль просто перезаписывает то же состояние (upsert по ключу события). Это дешевле и надёжнее, чем честный exactly-once через границы разных систем.

// Порядок коммита решает семантику потребителя: сначала обработать, потом коммитить offset это at-least-once (при падении переобработаем); коммитить offset до обработки - at-most-once (при падении потеряем).

at-least-once
без потерь, возможны дубли - дефолт индустрии
идемпотентный sink
запись, где повтор не меняет результат (upsert)

Откуда дубли

Дубли рождаются на ретраях с обеих сторон. Продюсер не получил ack и повторил отправку - сообщение записалось дважды. Потребитель обработал, но упал до коммита offset'а - после перезапуска обработает ту же пачку снова.

Идемпотентный продюсер Kafka (enable.idempotence) убирает дубли ретраев в пределах партиции. Транзакции продюсера дают атомарную запись сразу в несколько партиций плюс read_committed у читателя.

// Автокоммит по таймеру при долгой обработке особенно коварен: после ребаланса пачка переобработается, потому что offset закоммитился раньше, чем работа доделалась.

transactional producer
атомарная запись в несколько партиций Kafka
read_committed
читатель видит только закоммиченные транзакции

Exactly-once и его границы

Exactly-once processing (Kafka Streams, Flink) это атомарность тройки «состояние + offset + выдача» внутри пайплайна. Внутри одной поддерживающей системы оно реально.

Но на выходе во внешний мир - БД, API - гарантия обрывается: там всё равно нужен идемпотентный sink или транзакционная запись. Дедупликацию на границе делают ключом идемпотентности события (event_id) плюс окном хранения виденных ключей.

// Выбор семантики - по цене ошибки: биллинг не терпит дублей, там идемпотентность обязательна; для метрик потеря доли процента дешевле, чем инфраструктура exactly-once. Строить EOS (end of sequence) ради метрик - избыточно.

exactly-once
эффект ровно один раз в пределах поддерживающих систем
дедуп по event_id
ключ события + окно виденных ключей на границе

Как отвечать: «Бывает ли exactly-once на самом деле?»

В узких границах - да, но как только пересекаешь границу систем, честнее говорить про at-least-once плюс идемпотентность. Exactly-once processing реально внутри одного движка вроде Kafka Streams или Flink: они атомарно фиксируют тройку - своё состояние, прочитанный offset и выданный результат, поэтому в пределах пайплайна сбой не потеряет и не задвоит. Но когда пайплайн пишет наружу, в базу или в чужой API, эта атомарность обрывается: внешняя система про транзакцию движка ничего не знает. Поэтому на выходе всё равно нужен идемпотентный sink - upsert по ключу - или транзакционная запись, иначе повтор даст дубль. Отсюда практический вывод: я не полагаюсь на маркетинговое «у нас exactly-once», а проектирую потребителя идемпотентным, чтобы дубль перезаписывал то же состояние. А уровень усилий выбираю по цене ошибки: биллингу дубль недопустим, метрикам потеря 0.1% дешевле, чем EOS-инфраструктура.

Почему это сильный ответ: exactly-once не отрицается и не абсолютизируется - показана его реальная граница (внутри движка), необходимость идемпотентного sink наружу и выбор по цене ошибки.

На чём валят

  • Верить «у нас exactly-once», когда sink - внешняя БД без идемпотентности/транзакций.
  • Коммит offset'а до обработки - падение = молча потерянные события.
  • Автокоммит по таймеру + долгая обработка - переобработка пачки после ребаланса.
  • Дедуп по «всей истории» в памяти - состояние растёт бесконечно; нужно окно + TTL (time to live).
  • Строить exactly-once ради метрик, где дубль в 0.1% никого не волнует.

Проверьте себя

Пять вопросов из банка по этой подтеме. Всего их 16, остальные разбираются в тренажёре.

  1. #delivery_semantics1 / 5
    Что обеспечивает exactly-once semantics в Kafka?
    A)Обязательный полный отказ от партиционирования топиков в пользу одной общей очереди
    B)Идемпотентный продюсер + транзакции: атомарный read-process-write без дублей
    C)Простое увеличение числа реплик каждой партиции топика ровно втрое
    D)Ручное удаление дубликатов аналитиком раз в сутки по расписанию
    показать ответ и разбор
    +B)Идемпотентный продюсер + транзакции: атомарный read-process-write без дублей

    // разбор: EOS в Kafka строится на идемпотентном продюсере (брокер отсекает дубли по producer id + sequence) и транзакционном API, который делает цепочку «прочитал → обработал → записал + закоммитил offset» атомарной. Так consume-transform-produce не порождает дублей даже при сбоях. Партиционирование при этом сохраняется, а реплики отвечают за надёжность хранения, а не за дубли.

  2. #delivery_semantics2 / 5
    Почему хранение событий как неизменяемого лога позволяет «переиграть» (replay) обработку?
    A)Лог событий хранит одно самое последнее событие, перезаписывая все прошлые
    B)Replay не выходит: обработанные события удаляются из Kafka сразу после чтения и восстановить их оттуда для перечитывания уже не получится
    C)Переиграть можно, остановив всех текущих продюсеров
    D)События остаются в логе; можно перечитать их с нужного offset новой версией кода
    показать ответ и разбор
    +D)События остаются в логе; можно перечитать их с нужного offset новой версией кода

    // разбор: Kafka хранит сообщения как неизменяемый лог в пределах ретенции, а offset консьюмера — это просто указатель. Поэтому обработку можно переиграть: сбросить offset назад и прогнать те же события новой (исправленной) версией кода или в новый приёмник. Это основа event sourcing и способ восстановить/пересчитать витрину после багфикса без потери данных.

  3. #delivery_semantics3 / 5
    Что такое Change Data Capture (CDC)?
    A)Обязательная полная перезагрузка сразу всей таблицы целиком при каждом изменении даже одной строки
    B)Полная перезагрузка всей таблицы целиком по расписанию строго раз в сутки ночью
    C)Захват изменений строк БД (insert/update/delete) как потока событий
    D)Шифрование данных при их передаче между базой и хранилищем данных
    показать ответ и разбор
    +C)Захват изменений строк БД (insert/update/delete) как потока событий

    // разбор: CDC читает журнал изменений базы (WAL/binlog) и превращает каждый insert/update/delete в событие потока. Это позволяет реплицировать изменения в хранилище или другой сервис почти в реальном времени, без тяжёлых периодических полных выгрузок таблицы и без нагрузки тяжёлыми SELECT на источник.

  4. #delivery_semantics4 / 5
    Что задаёт параметр acks у продюсера Kafka?
    A)Число сообщений, которое продюсер накапливает в батч перед отправкой брокеру Kafka
    B)Сколько реплик должны подтвердить запись, прежде чем продюсер считает её успешной (0/1/all)
    C)Максимальное число консьюмеров, которым брокер разрешит одновременно читать данный топик
    D)Сколько раз продюсер повторит отправку сообщения при возникновении сетевой ошибки
    показать ответ и разбор
    +B)Сколько реплик должны подтвердить запись, прежде чем продюсер считает её успешной (0/1/all)

    // разбор: acks — гарантия долговечности записи. acks=0: продюсер не ждёт подтверждения (быстро, можно потерять). acks=1: ждёт лидера партиции (потеря при падении лидера до репликации). acks=all: ждёт все синхронные реплики (ISR) — не теряется, пока жива хоть одна реплика, но выше задержка. Выбор acks — размен между надёжностью и скоростью; для важных данных берут all.

  5. #delivery_semantics5 / 5
    Что такое offset в Kafka и зачем консьюмер его коммитит?
    A)Offset — это физический адрес сообщения на диске брокера, меняющийся при каждом слиянии сегментов
    B)Offset — приоритет сообщения, определяющий, в каком порядке консьюмеры получат его из топика
    C)Порядковый номер сообщения в партиции; коммит offset помечает, докуда прочитано, чтобы продолжить после перезапуска
    D)Offset уникален глобально по всему топику сразу, а не в пределах отдельной партиции
    показать ответ и разбор
    +C)Порядковый номер сообщения в партиции; коммит offset помечает, докуда прочитано, чтобы продолжить после перезапуска

    // разбор: Offset — монотонный порядковый номер записи внутри партиции. Консьюмер читает по возрастанию offset и периодически коммитит достигнутую позицию (в служебный топик), чтобы после перезапуска/ребаланса продолжить с неё, а не с начала. Момент коммита определяет семантику: коммит до обработки — риск потери (at-most-once), после — риск повтора при сбое (at-least-once). Отсюда и выбор гарантий.

дальше

Теорию прочитали. Навык ставится повторением

В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.