Гарантии доставки в 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, остальные разбираются в тренажёре.
- Что обеспечивает 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 не порождает дублей даже при сбоях. Партиционирование при этом сохраняется, а реплики отвечают за надёжность хранения, а не за дубли.
- Почему хранение событий как неизменяемого лога позволяет «переиграть» (replay) обработку?A)Лог событий хранит одно самое последнее событие, перезаписывая все прошлыеB)Replay не выходит: обработанные события удаляются из Kafka сразу после чтения и восстановить их оттуда для перечитывания уже не получитсяC)Переиграть можно, остановив всех текущих продюсеровD)События остаются в логе; можно перечитать их с нужного offset новой версией кода
показать ответ и разбор
+D)События остаются в логе; можно перечитать их с нужного offset новой версией кода// разбор: Kafka хранит сообщения как неизменяемый лог в пределах ретенции, а offset консьюмера — это просто указатель. Поэтому обработку можно переиграть: сбросить offset назад и прогнать те же события новой (исправленной) версией кода или в новый приёмник. Это основа event sourcing и способ восстановить/пересчитать витрину после багфикса без потери данных.
- Что такое Change Data Capture (CDC)?A)Обязательная полная перезагрузка сразу всей таблицы целиком при каждом изменении даже одной строкиB)Полная перезагрузка всей таблицы целиком по расписанию строго раз в сутки ночьюC)Захват изменений строк БД (insert/update/delete) как потока событийD)Шифрование данных при их передаче между базой и хранилищем данных
показать ответ и разбор
+C)Захват изменений строк БД (insert/update/delete) как потока событий// разбор: CDC читает журнал изменений базы (WAL/binlog) и превращает каждый insert/update/delete в событие потока. Это позволяет реплицировать изменения в хранилище или другой сервис почти в реальном времени, без тяжёлых периодических полных выгрузок таблицы и без нагрузки тяжёлыми SELECT на источник.
- Что задаёт параметр acks у продюсера Kafka?A)Число сообщений, которое продюсер накапливает в батч перед отправкой брокеру KafkaB)Сколько реплик должны подтвердить запись, прежде чем продюсер считает её успешной (0/1/all)C)Максимальное число консьюмеров, которым брокер разрешит одновременно читать данный топикD)Сколько раз продюсер повторит отправку сообщения при возникновении сетевой ошибки
показать ответ и разбор
+B)Сколько реплик должны подтвердить запись, прежде чем продюсер считает её успешной (0/1/all)// разбор: acks — гарантия долговечности записи. acks=0: продюсер не ждёт подтверждения (быстро, можно потерять). acks=1: ждёт лидера партиции (потеря при падении лидера до репликации). acks=all: ждёт все синхронные реплики (ISR) — не теряется, пока жива хоть одна реплика, но выше задержка. Выбор acks — размен между надёжностью и скоростью; для важных данных берут all.
- Что такое offset в Kafka и зачем консьюмер его коммитит?A)Offset — это физический адрес сообщения на диске брокера, меняющийся при каждом слиянии сегментовB)Offset — приоритет сообщения, определяющий, в каком порядке консьюмеры получат его из топикаC)Порядковый номер сообщения в партиции; коммит offset помечает, докуда прочитано, чтобы продолжить после перезапускаD)Offset уникален глобально по всему топику сразу, а не в пределах отдельной партиции
показать ответ и разбор
+C)Порядковый номер сообщения в партиции; коммит offset помечает, докуда прочитано, чтобы продолжить после перезапуска// разбор: Offset — монотонный порядковый номер записи внутри партиции. Консьюмер читает по возрастанию offset и периодически коммитит достигнутую позицию (в служебный топик), чтобы после перезапуска/ребаланса продолжить с неё, а не с начала. Момент коммита определяет семантику: коммит до обработки — риск потери (at-most-once), после — риск повтора при сбое (at-least-once). Отсюда и выбор гарантий.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.