Messaging на JVM: Kafka-клиенты, Spring Kafka и надёжная доставка

Работа с Kafka из JVM: официальный клиент, Spring Kafka, гарантии доставки и идемпотентность, сериализация и schema registry — что важно кроме «отправить сообщение»

Отправить сообщение в 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-клиенте и надёжности.

Kafka на JVM: consumer groups, exactly-once и dead-letter

В статье

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, получено=20

20 отправлено, 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

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
Ребаланс: consumer-2 подключается, партиции переназначаются без пересечений

Результат честный: 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

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
EOS: абортнутый батч физически в логе, но невидим read_committed-консьюмеру

Здесь важно физическое и логическое число не путать: 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

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
Retry с backoff, затем DLT + commitSync

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».

Документация и первоисточники

Обсуждение в Telegram

Присоединиться →

Комментарии