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

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=all
  • replication.factor=3
  • min.insync.replicas=2

Retries и дубликаты

Если producer не получил ответ, он не всегда понимает:

  • запись не дошла вовсе;
  • запись дошла, но ack потерялся;
  • запись записалась частично;
  • metadata устарела и лидер уже сменился.

Поэтому producer делает retry.

Но retry сам по себе открывает риск дублей:

send -> network error -> 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 вида:

consume from topic A
process
produce to topic B
commit offsets

Если сделать это неаккуратно, можно получить:

  • output записан, а offset не закоммичен;
  • offset закоммичен, а output не записан;
  • дубликаты при ретраях и рестартах.

Kafka transactions помогают решить эту проблему внутри Kafka pipeline.


Consumer: как работает чтение

Consumer обычно работает так:

  1. входит в consumer group;
  2. подписывается на topic;
  3. получает assignment partition'ов;
  4. вызывает poll;
  5. получает batch records;
  6. обрабатывает их;
  7. коммитит offset.

Короткая схема:

subscribe -> poll -> process -> commit -> poll -> process -> commit

Позиция 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 до обработки

poll -> commit -> process

Если приложение упало после commit, но до обработки:

  • сообщения потеряны для этой group.

Это at-most-once стиль.

Commit после обработки

poll -> process -> commit

Если приложение обработало запись, но упало до commit:

  • после рестарта запись придет снова.

Это at-least-once стиль.

На практике именно этот стиль чаще всего и используется.


Consumer group и распределение partition

Для одного topic внутри одной group:

  • одна partition в стандартном consumer group назначается одному consumer;
  • один consumer может обслуживать несколько partition;
  • число реально активных consumer'ов ограничено количеством partition.

Пример:

Topic has 4 partitions

Group:
  consumer-1 -> p0, p1
  consumer-2 -> p2
  consumer-3 -> p3

Если в group добавить consumer-4, он получит работу только если есть свободная partition после rebalance.


Rebalance

Rebalance — это перераспределение partition между consumer'ами группы.

Он может происходить, когда:

  • consumer присоединился к group;
  • consumer вышел из group;
  • consumer перестал слать heartbeats;
  • изменилась подписка;
  • изменилось число partition у topic.

Проблема rebalance в том, что:

  • assignment меняется;
  • consumer может потерять partition прямо во время обработки;
  • если код плохо написан, легко словить дубли или пропуски.

Поэтому rebalance — это не мелкая техническая деталь, а важная часть дизайна consumer'а.


Почему rebalance часто ломает наивную обработку

Типичный плохой сценарий:

  1. Consumer получил batch.
  2. Начал долгую обработку.
  3. Перестал успевать в heartbeat / poll cadence.
  4. Kafka решила, что consumer завис.
  5. Произошел rebalance.
  6. Partition отдали другому consumer.
  7. Часть записей может быть обработана повторно.

Поэтому:

  • не надо держать слишком долгую обработку внутри одного 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.