Kafka Producers & Consumers¶
Producer: что делает producer¶
Producer:
- сериализует событие;
- выбирает topic;
- выбирает partition;
- отправляет record лидеру partition;
- ждет ack в зависимости от настроек;
- при необходимости делает retry;
- может использовать batching и compression.
Producer обычно является началом всей data flow.
Как producer выбирает partition¶
Есть три основных варианта:
1. Явно указать partition¶
Редко нужно вручную, но иногда полезно для специальных сценариев.
Минусы:
- сильная привязка к текущей структуре topic;
- легко получить перекос нагрузки.
2. Передать key¶
Самый частый вариант.
Тогда одна и та же key обычно хешируется в одну и ту же partition.
Например:
- все события по
orderId=123-> в одну partition; - все события по
userId=42-> в одну partition.
Это дает:
- локальный порядок событий по одной сущности;
- предсказуемое распределение данных.
3. Не передавать key¶
Тогда producer использует встроенную стратегию распределения.
Идея:
- балансировать записи по partition;
- максимизировать throughput;
- но не гарантировать порядок для одной конкретной сущности.
Практическое правило:
- есть сущность, для которой нужен порядок -> обязательно думай про key;
- нет такой сущности -> можно писать без key ради более равномерной нагрузки.
Batching и compression¶
Kafka producer эффективен потому, что умеет отправлять данные пачками.
Producer:
- буферизует записи;
- собирает batch;
- отправляет больше данных одним запросом.
Это дает:
- меньше сетевых round-trip;
- выше throughput;
- лучше компрессию.
Компромисс:
- больше batching -> выше throughput;
- но чуть выше latency.
На практике для Kafka это нормальная tradeoff-модель.
acks: сколько подтверждений ждать¶
acks определяет, когда producer считает запись успешной.
acks=0¶
Producer вообще не ждет подтверждения.
Плюсы:
- минимальная latency;
- максимальный throughput.
Минусы:
- producer не знает, дошла запись или нет;
- высокий риск silent loss.
acks=1¶
Producer ждет подтверждение только от leader.
Плюсы:
- быстрее, чем
acks=all; - популярный компромисс для некритичных потоков.
Минусы:
- если leader подтвердил запись, но умер до репликации, запись может потеряться.
acks=all или acks=-1¶
Producer ждет подтверждение на уровне текущего ISR.
Плюсы:
- strongest durability из стандартных producer mode;
- хороший выбор для бизнес-критичных событий.
Важно:
- реальная гарантия зависит не только от
acks=all; - она зависит еще от
replication.factorиmin.insync.replicas.
Типичный надежный вариант:
acks=allreplication.factor=3min.insync.replicas=2
Retries и дубликаты¶
Если producer не получил ответ, он не всегда понимает:
- запись не дошла вовсе;
- запись дошла, но ack потерялся;
- запись записалась частично;
- metadata устарела и лидер уже сменился.
Поэтому producer делает retry.
Но retry сам по себе открывает риск дублей:
Вывод:
- retries полезны;
- но без дополнительных механизмов они могут давать duplicate records.
Idempotent producer¶
Kafka поддерживает idempotent producer.
Идея:
- broker выдает producer ID;
- producer отправляет sequence numbers;
- broker умеет отличать повторную отправку того же record batch.
Это уменьшает риск дублей из-за retry на уровне записи в Kafka log.
Важно понимать:
- idempotent producer защищает именно запись в Kafka;
- он не делает магически идемпотентной вашу бизнес-логику;
- если consumer дважды вызвал внешний API, Kafka это уже не исправит.
Producer transactions¶
Kafka поддерживает transactions.
Они нужны, когда надо атомарно:
- записать данные в несколько topic/partition;
- и/или
- записать output records вместе с offset commit в Kafka.
Это важно для pipeline вида:
Если сделать это неаккуратно, можно получить:
- output записан, а offset не закоммичен;
- offset закоммичен, а output не записан;
- дубликаты при ретраях и рестартах.
Kafka transactions помогают решить эту проблему внутри Kafka pipeline.
Consumer: как работает чтение¶
Consumer обычно работает так:
- входит в consumer group;
- подписывается на topic;
- получает assignment partition'ов;
- вызывает
poll; - получает batch records;
- обрабатывает их;
- коммитит offset.
Короткая схема:
Позиция consumer: current position vs committed offset¶
Нужно различать:
- current position — что consumer уже прочитал в памяти прямо сейчас;
- committed offset — до какого места группа считает обработку безопасно завершенной.
Это не одно и то же.
Consumer может:
- уже получить записи через
poll; - начать их обрабатывать;
- но еще не успеть commit offset.
Если в этот момент consumer упадет, после рестарта он начнет от committed offset, а значит часть записей может прийти повторно.
Auto commit и manual commit¶
Auto commit¶
Kafka может автоматически периодически коммитить offsets.
Плюсы:
- проще стартовать;
- меньше кода.
Минусы:
- легко случайно закоммитить offset раньше, чем запись реально обработана;
- хуже контроль при сложной бизнес-логике.
Manual commit¶
Приложение само решает, когда коммитить offset.
Плюсы:
- полный контроль;
- лучше подходит для критичных сценариев.
Минусы:
- больше ответственности;
- легче ошибиться в коде и создать дубли или потери.
Практически:
- для серьезных consumer'ов обычно нужен осознанный manual commit;
- auto commit хорошо подходит только для простых или учебных сценариев.
Когда коммитить offset¶
Это один из самых важных вопросов.
Commit до обработки¶
Если приложение упало после commit, но до обработки:
- сообщения потеряны для этой group.
Это at-most-once стиль.
Commit после обработки¶
Если приложение обработало запись, но упало до commit:
- после рестарта запись придет снова.
Это at-least-once стиль.
На практике именно этот стиль чаще всего и используется.
Consumer group и распределение partition¶
Для одного topic внутри одной group:
- одна partition в стандартном consumer group назначается одному consumer;
- один consumer может обслуживать несколько partition;
- число реально активных consumer'ов ограничено количеством partition.
Пример:
Если в group добавить consumer-4, он получит работу только если есть свободная partition после rebalance.
Rebalance¶
Rebalance — это перераспределение partition между consumer'ами группы.
Он может происходить, когда:
- consumer присоединился к group;
- consumer вышел из group;
- consumer перестал слать heartbeats;
- изменилась подписка;
- изменилось число partition у topic.
Проблема rebalance в том, что:
- assignment меняется;
- consumer может потерять partition прямо во время обработки;
- если код плохо написан, легко словить дубли или пропуски.
Поэтому rebalance — это не мелкая техническая деталь, а важная часть дизайна consumer'а.
Почему rebalance часто ломает наивную обработку¶
Типичный плохой сценарий:
- Consumer получил batch.
- Начал долгую обработку.
- Перестал успевать в heartbeat / poll cadence.
- Kafka решила, что consumer завис.
- Произошел rebalance.
- Partition отдали другому consumer.
- Часть записей может быть обработана повторно.
Поэтому:
- не надо держать слишком долгую обработку внутри одного poll cycle без контроля;
- нужно понимать
max.poll.interval.ms, batch size и время обработки; - idempotent business processing очень желательна.
Delivery semantics¶
At-most-once¶
Смысл:
- запись обрабатывается не более одного раза;
- но может потеряться.
Как обычно получается:
- commit offset раньше фактической обработки.
At-least-once¶
Смысл:
- запись не теряется;
- но может быть обработана повторно.
Как обычно получается:
- сначала обработка;
- потом commit.
Это самый частый реальный режим.
Effectively-once¶
Это не официальный режим протокола, а практический результат:
- Kafka может прислать запись повторно;
- но приложение делает эффект идемпотентно;
- поэтому с точки зрения бизнеса дубликат безвреден.
Пример:
- upsert по
orderId; - запись в БД с уникальным constraint;
- deduplication table по event ID.
Exactly-once¶
Тут важно не переоценивать термин.
В Kafka exactly-once означает ограниченный сценарий:
- consumer читает из Kafka;
- producer пишет обратно в Kafka;
- offsets и output records коммитятся транзакционно;
- downstream читает
read_committed.
То есть:
- это не "абсолютно без дублей во всей системе";
- это не "внешняя БД и внешний API автоматически стали exactly-once";
- это exactly-once в рамках корректно построенного Kafka-to-Kafka pipeline.
read_committed¶
Если используются Kafka transactions, downstream consumer обычно должен читать в режиме read_committed.
Иначе он может увидеть:
- записи из abort'нутых транзакций;
- промежуточные данные, которые не должны считаться финальными.
Это важно для реального exactly-once pipeline.
Практические советы по consumer'ам¶
- Думай об offset commit как о boundary между "можно повторить" и "считаем завершенным".
- Закладывайся на дубликаты почти всегда.
- Делай обработку идемпотентной.
- Не держи слишком долгую обработку без понимания rebalance behavior.
- Учитывай, что увеличение partition влияет на parallelism, но и на сложность тоже.
- Не включай auto commit без понимания, что именно ты теряешь в контроле.
Что запомнить¶
- Producer выбирает topic, key, partition и стратегию подтверждения.
acks=allбез нормальногоmin.insync.replicas— это еще не вся durability story.- Retry без idempotence может создать дубликаты.
- Consumer управляет offset, а не "удаляет сообщения".
- Rebalance — нормальная часть жизни consumer group, а не редкая авария.
- Самый частый реальный режим обработки — at-least-once + идемпотентная бизнес-логика.
- Exactly-once в Kafka имеет смысл только в четко ограниченном pipeline.