Отправить сообщение в Kafka — одна строка. Не потерять его, не задублировать и корректно обработать сбой consumer’а — совсем другая история. JVM — родная платформа для Kafka: kafka-clients — референсная реализация протокола, и Spring Kafka, и большинство JVM-фреймворков строятся поверх неё же.
Здесь не про Kafka целиком — про ментальную модель лога, партиций и retention в серии messaging есть отдельная статья. Фокус этой статьи — то, что видно именно со стороны JVM-клиента: как реально ведут себя consumer groups при ребалансе, что даёт exactly-once semantics не на бумаге, а в конкретном прогоне, как обрабатывать ошибки без потери сообщений и блокировки партиции, и где в этой картине место Spring Kafka и Kotlin-корутин.
Все примеры и цифры в статье — из живого стенда digital-cookbook/java/deep-dive/messaging: реальный брокер apache/kafka:4.3.0 в режиме KRaft (single-node), сырой kafka-clients без Spring, четыре сценария, каждый — с ассертом на итоговые числа, а не с придуманным выводом.
Статья JVM-серии; пересекается с брокерами сообщений на уровне выбора инструмента, но здесь фокус на JVM-клиенте и надёжности.
В статье
- Producer, consumer и живой прогон
- Consumer groups и ребаланс
- Exactly-once: идемпотентность и транзакции
- Обработка ошибок: retry и dead-letter topic
- Spring Kafka и Kotlin-корутины
- Сериализация и Schema Registry
- Когда это действительно важно
Producer, consumer и живой прогон
Базовая модель kafka-clients не сложная: KafkaProducer пишет записи в топик (сериализация ключа/значения, партиционирование по ключу или round-robin), KafkaConsumer читает их через poll(), продвигая offset — позицию последнего прочитанного сообщения в партиции. Consumer group — это способ разделить чтение топика между несколькими экземплярами consumer’а: каждая партиция читается ровно одним consumer’ом внутри группы (подробнее про сам лог и партиции — в «Kafka: лог, топики, партиции»).
На стенде это выглядит так: топик demo-basic с 3 партициями, producer отправляет 20 сообщений с ключами key-0/key-1/key-2 (партиционер выбирает партицию детерминированно по хешу ключа — одинаковый ключ всегда попадает в ту же партицию, что даёт упорядоченность в пределах ключа, но разные ключи не гарантированно оказываются в разных партициях — возможны хеш-коллизии), consumer с auto.offset.reset=earliest и enable.auto.commit=true читает их обратно.
Properties producerProps = new Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Kafka.BOOTSTRAP);
producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
producerProps.put(ProducerConfig.ACKS_CONFIG, "all");
int sent = 0;
try (KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps)) {
for (int i = 0; i < messageCount; i++) {
String key = "key-" + (i % PARTITIONS);
String value = "msg-" + i;
RecordMetadata meta = producer.send(new ProducerRecord<>(TOPIC, key, value)).get();
sent++;
}
}Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, Kafka.BOOTSTRAP);
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "basic-group");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
int received = 0;
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps)) {
consumer.subscribe(List.of(TOPIC));
while (received < messageCount) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
received++;
}
}
}[producer] всего отправлено: 20
[consumer] всего получено: 20
АССЕРТ OK: отправлено=20, получено=2020 отправлено, 20 получено — ожидаемо для happy path с acks=all. Обратите внимание на enable.auto.commit=true: consumer коммитит offset периодически в фоне, независимо от того, успешно ли приложение обработало сообщение. Это простая настройка для demo, но именно она — источник большинства «пропаданий» и дублей сообщений в реальных системах; к этому мы вернёмся в разделе про exactly-once и dead-letter topic.
Ещё одна деталь, которая всплывает при первом запуске: если топик создан через AdminClient прямо перед отправкой в него, KafkaProducer иногда получает NotLeaderOrFollowerException на первой попытке — метаданные о свежесозданном топике ещё не разошлись между клиентом и брокером. Это штатное поведение single-node брокера (задержка в доли секунды), а не ошибка: KafkaProducer автоматически ретраит с обновлением метаданных, в логе виден WARN, сообщения доставляются. В проде с этим сталкиваются реже — топики обычно уже существуют к моменту продьюсинга, — но сам механизм авто-ретрая универсален для любого кластера.
Consumer groups и ребаланс
Самое интересное в consumer groups — не стационарное состояние, а момент, когда состав группы меняется: подключается новый consumer, отваливается старый, меняется число партиций. Именно тогда происходит ребаланс — брокер (group coordinator) отзывает партиции у текущих участников и переназначает их заново.
На стенде это видно напрямую. Топик demo-groups с 4 партициями засеян 40 сообщениями заранее. Сначала подключается consumer-1 — единственный участник группы, он получает все 4 партиции. Затем подключается consumer-2 — это триггерит ребаланс, который ловится через ConsumerRebalanceListener:
consumer.subscribe(List.of(TOPIC), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
if (!partitions.isEmpty()) {
System.out.printf(" [rebalance] %s: revoked %s%n", id, partitions);
}
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
lastAssignment = Set.copyOf(partitions);
System.out.printf(" [rebalance] %s: assigned %s%n", id, partitions);
firstAssignment.countDown();
}
});[rebalance] consumer-1: assigned [demo-groups-0, demo-groups-1, demo-groups-2, demo-groups-3]
[rebalance] consumer-1: revoked [demo-groups-0, demo-groups-1, demo-groups-2, demo-groups-3]
[rebalance] consumer-1: assigned [demo-groups-0, demo-groups-1]
[rebalance] consumer-2: assigned [demo-groups-2, demo-groups-3]
АССЕРТ OK: партиции распределены без пересечений (2+2=4 из 4), сообщений дочитано=40 (>=40)
flowchart LR
subgraph before ["До подключения consumer-2"]
C1a["consumer-1"] --- P0["partition-0"]
C1a --- P1["partition-1"]
C1a --- P2["partition-2"]
C1a --- P3["partition-3"]
end
subgraph after ["После ребаланса"]
C1b["consumer-1"] --- Q0["partition-0"]
C1b --- Q1["partition-1"]
C2b["consumer-2"] --- Q2["partition-2"]
C2b --- Q3["partition-3"]
end
before -->|"consumer-2 подключается\nrevoke + assign"| after
style C1a fill:#f9f3e3,stroke:#8b7355
style C1b fill:#c9e4c5,stroke:#5b8a5e
style C2b fill:#c9e4c5,stroke:#5b8a5e
Результат честный: 2+2 партиции без пересечений, onPartitionsRevoked у consumer-1 перед новым onPartitionsAssigned у обоих — ровно то поведение, которое описывает протокол consumer group. Но есть нюанс в самом сценарии, который стоит проговорить прямо: топик был засеян сообщениями до подключения консьюмеров, поэтому consumer-1, будучи какое-то время единственным участником группы, успел дочитать почти всё ещё до ребаланса (consumer-1 дочитал 40, consumer-2 — 0). Партиционирование при этом распределилось корректно, но «живое» разделение чтения между двумя работающими consumer’ами на этом прогоне наглядно не видно. Для демонстрационных целей нагляднее обратный порядок — сначала поднять обоих consumer’ов, затем продюсить непрерывно, пока оба работают: тогда видно, как каждый читает именно свою часть потока, а не то, что один успел вычерпать топик в одиночку.
Подробнее про сам механизм group coordinator, partition assignor’ы и статичное членство группы (static membership) — в «Consumer groups и ребалансировка в Kafka».
Exactly-once: идемпотентность и транзакции
У Kafka три уровня гарантий доставки: at-most-once (можно потерять, дублей не будет — обычно enable.auto.commit с коммитом до обработки), at-least-once (дубли возможны, потерь нет — коммит после обработки) и exactly-once semantics (EOS) — единственное сообщение появляется у consumer’а ровно один раз, даже если producer или обработка падали и ретраились.
EOS в Kafka строится на двух механизмах: идемпотентном producer’е (enable.idempotence=true — брокер дедуплицирует повторные записи по producer id + sequence number на уровне партиции) и транзакциях (transactional.id — атомарная запись в несколько партиций/топиков, видимая consumer’у с isolation.level=read_committed только целиком, при commitTransaction(), и полностью невидимая при abortTransaction()).
На стенде — топик demo-eos, транзакционный producer с transactional.id=messaging-eos-producer и enable.idempotence=true. Три батча по 5 сообщений: батч A коммитится штатно; батч B на первой попытке «падает» в середине (симуляция сбоя обработки) и абортится; батч B пересылается повторно (тот же контент) и коммитится штатно.
private static void sendBatch(KafkaProducer<String, String> producer, String prefix, boolean failMidway) {
producer.beginTransaction();
try {
for (int i = 0; i < BATCH_SIZE; i++) {
producer.send(new ProducerRecord<>(TOPIC, prefix + "-key-" + i, prefix + "-" + i));
if (failMidway && i == 2) {
throw new RuntimeException("симулированный сбой обработки в середине батча " + prefix);
}
}
producer.commitTransaction();
System.out.printf(" [eos] батч %s: commitTransaction() — %d сообщений подтверждено%n", prefix, BATCH_SIZE);
} catch (RuntimeException e) {
System.out.printf(" [eos] батч %s: сбой (%s) -> abortTransaction()%n", prefix, e.getMessage());
producer.abortTransaction();
}
}producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "messaging-eos-producer");
producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
producerProps.put(ProducerConfig.ACKS_CONFIG, "all");
// ...
producer.initTransactions();
sendBatch(producer, "batchA", false); // штатный коммит
sendBatch(producer, "batchB", true); // сбой в середине -> abort
sendBatch(producer, "batchB", false); // повтор того же батча -> коммит
// consumer:
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");физически отправлено записей: 15 (включая абортнутый батч)
логически подтверждено: 10
увидел read_committed-консьюмер: 10
АССЕРТ OK: read_committed увидел=10 (абортнутый батч невидим, дублей нет)
flowchart LR
A["batch A\n5 msg"] -->|"commitTransaction()"| V1["видно read_committed"]
B1["batch B, попытка 1\n5 msg"] -->|"сбой в середине\nabortTransaction()"| X["невидно read_committed\ncontrol-record ABORT"]
B2["batch B, попытка 2\n5 msg"] -->|"commitTransaction()"| V2["видно read_committed"]
style X fill:#f5d6d0,stroke:#b05050
style V1 fill:#c9e4c5,stroke:#5b8a5e
style V2 fill:#c9e4c5,stroke:#5b8a5e
Здесь важно физическое и логическое число не путать: 15 записей реально лежат в логе партиций (Kafka не переписывает лог задним числом), но абортнутый батч помечен control-record’ом ABORT, и read_committed-consumer его полностью отфильтровывает. Итог — ровно 10, без дублей от повторной отправки того же батча B. Это и есть разница между «данные физически есть» и «данные логически подтверждены» — та же идея, что в WAL и MVCC других СУБД, только на уровне лога партиций Kafka.
Подробнее про сам протокол транзакций (двухфазный коммит через transaction coordinator, __transaction_state топик, producer epoch и fencing) — в «Exactly-once и транзакции в Kafka».
Обработка ошибок: retry и dead-letter topic
EOS решает проблему дублей и потерь на уровне транспорта, но не решает проблему «а что если сообщение в принципе не обрабатывается» — битый формат, недоступный внешний сервис, бизнес-ошибка. Здесь работает другой паттерн: несколько попыток с backoff, а после исчерпания ретраев — перенос в dead-letter topic (DLT), чтобы не блокировать партицию навсегда одним «ядовитым» сообщением.
На стенде — топик demo-errors (3 партиции) + demo-errors-dlt (1 партиция), 15 сообщений, каждое 5-е помечено маркером-«ядом» (обработка которого всегда падает). На каждое ядовитое сообщение — 3 попытки с backoff 150мс × номер попытки, затем отправка в DLT с заголовками об исходном топике/партиции/оффсете/причине и commitSync():
for (ConsumerRecord<String, String> r : records) {
boolean ok = false;
Exception lastError = null;
for (int attempt = 1; attempt <= MAX_RETRIES && !ok; attempt++) {
try {
process(r.value());
ok = true;
} catch (RuntimeException e) {
lastError = e;
if (attempt < MAX_RETRIES) {
Kafka.sleep(150L * attempt); // линейный бэкофф: 150 мс × номер попытки
}
}
}
if (ok) {
result.goodProcessed++;
} else {
List<Header> headers = List.of(
new RecordHeader("x-original-topic", TOPIC.getBytes(StandardCharsets.UTF_8)),
new RecordHeader("x-original-partition",
String.valueOf(r.partition()).getBytes(StandardCharsets.UTF_8)),
new RecordHeader("x-original-offset",
String.valueOf(r.offset()).getBytes(StandardCharsets.UTF_8)),
new RecordHeader("x-error",
String.valueOf(lastError).getBytes(StandardCharsets.UTF_8))
);
dltProducer.send(new ProducerRecord<>(DLT_TOPIC, null, r.key(), r.value(), headers));
dltProducer.flush();
result.sentToDlt++;
}
// Оффсет коммитится в любом случае (успех или отправка в DLT) —
// "ядовитое" сообщение не блокирует партицию навсегда.
consumer.commitSync();
processedTotal++;
}Структура заголовков, которые producer пишет в DLT-запись при исчерпании ретраев (значения подставляются из конкретного упавшего сообщения, здесь — схема, а не дословный вывод одного прогона):
x-original-topic: demo-errors
x-original-partition: <партиция исходного сообщения>
x-original-offset: <offset исходного сообщения>
x-error: <текст исключения, из-за которого исчерпаны 3 попытки>итог: обработано успешно=12, в DLT=3 (перепроверено чтением DLT=3)
АССЕРТ OK: успешно обработано=12, в DLT после 3 ретраев=3 (перепроверено=3)
flowchart TD
R["consumer.poll()"] --> A1["попытка 1/3"]
A1 -->|"OK"| C["commitSync()\ngoodProcessed++"]
A1 -->|"ошибка → backoff 150мс×1"| A2["попытка 2/3"]
A2 -->|"OK"| C
A2 -->|"ошибка → backoff 150мс×2"| A3["попытка 3/3"]
A3 -->|"OK"| C
A3 -->|"ошибка, ретраи исчерпаны"| D["DLT: заголовки x-original-*\n+ commitSync()"]
style D fill:#f5d6d0,stroke:#b05050
style C fill:#c9e4c5,stroke:#5b8a5e
12 успешно обработано, 3 в DLT, независимая перепроверка чтением DLT-топика тоже даёт 3 — сходится по всем трём каналам подсчёта. Из трёх ядовитых сообщений ни одно не осталось необработанным и не заблокировало партицию — commitSync() вызывается и при успехе, и при финальной отправке в DLT.
Обратите внимание: здесь commitSync() вызывается на каждую запись отдельно (per-record commit) — это удобно для демонстрации (видно, что оффсет продвигается ровно после решения по конкретному сообщению), но в проде почти всегда эффективнее батчевый коммит — раз в N сообщений или раз в интервал времени, а не синхронный round-trip к брокеру на каждую запись. Per-record commitSync() — правильный выбор только там, где сама стоимость лишнего round-trip незначима на фоне остальной обработки.
Spring Kafka и Kotlin-корутины
Всё выше — сырой kafka-clients, намеренно: гарантии доставки, ребаланс и транзакции нагляднее видны без слоя абстракции поверх них. Spring Kafka не меняет эти гарантии, а упаковывает те же примитивы в более декларативный API: KafkaTemplate вместо ручного управления KafkaProducer, @KafkaListener вместо цикла poll(), ContainerFactory с настраиваемым CommonErrorHandler для retry/DLT (DefaultErrorHandler с DeadLetterPublishingRecoverer реализует практически тот же паттерн, что и код выше, но конфигурацией, а не ручным циклом), KafkaTransactionManager для интеграции транзакций Kafka с @Transactional. Понимание того, что происходит под капотом — то, что показано в этой статье, — облегчает чтение логов и диагностику, когда декларативный слой Spring Kafka скрывает детали.
Для Kotlin на JVM ситуация похожая: kafka-clients — блокирующий API, и корутины сюда добавляются либо через reactor-kafka (реактивный клиент, продюсер/консьюмер как Flux/Mono), либо через собственные suspend-обёртки поверх блокирующих вызовов, вынесенных в Dispatchers.IO. Ни то ни другое не меняет модель гарантий Kafka — она определяется конфигурацией producer/consumer, а не тем, как вызывающий код ждёт результата. Про сами идиомы Kotlin для backend, включая корутины, — в «Kotlin для backend».
Если источник событий в топике — не ваш сервис, а изменения в базе данных (change data capture), тот же consumer-side код (ребаланс, offset’ы, retry) применим и там; специфика CDC — в «CDC через Debezium: PostgreSQL → Kafka».
Сериализация и Schema Registry
Все примеры выше используют StringSerializer/StringDeserializer — сознательное упрощение для демонстрации гарантий доставки, а не сериализации. В проде выбор формата сообщений — отдельная развилка: JSON читаем и прост в отладке, но не проверяет схему на уровне брокера; Avro и Protobuf компактнее и типизированы, а в связке со Schema Registry дают проверку совместимости схемы при публикации (backward/forward compatibility) — consumer со старой версией схемы не падает на новом поле, если совместимость не нарушена намеренно.
Сама эволюция схемы без поломки уже работающих consumer’ов — тот же класс задачи, что версионирование контрактов между сервисами в целом, независимо от транспорта; см. «Обработка ошибок и контракты между языками». Практика работы со Schema Registry, Kafka Connect и ksqlDB — в «Экосистема Kafka: Connect, Schema Registry, ksqlDB».
Когда это действительно важно
Не каждому consumer’у нужен EOS и DLT. Простая аналитика, где потеря или дубль одного события — шум на фоне агрегата, отлично живёт на at-least-once с идемпотентной обработкой на стороне потребителя (например, upsert по бизнес-ключу вместо insert). Транзакционный producer и read_committed стоят своей сложности там, где дубль или потеря — это не шум, а некорректная сумма на счету, задвоенный платёж или рассинхронизация состояния между сервисами.
DLT и retry с backoff — более универсальный паттерн: почти любой consumer в проде должен уметь не падать намертво на одном плохом сообщении и не блокировать партицию навсегда. Тестировать эту логику — в том числе против настоящего брокера, а не мока — с той же связкой JUnit 5 + Testcontainers, что и для баз данных; см. «Тестирование на JVM».
Документация и первоисточники
- Стенд:
digital-cookbook/java/deep-dive/messaging(живойapache/kafka:4.3.0в KRaft, сыройkafka-clients; четыре сценария — базовая доставка, ребаланс consumer group, exactly-once через транзакции, retry + dead-letter topic — каждый с ассертом на итоговые числа) - Apache Kafka Documentation
- Kafka: Exactly-Once Semantics
- Spring for Apache Kafka
- Confluent Schema Registry
Комментарии