На стенде Kafka не подводит: продьюсер пишет, слушатель читает. Первая неожиданность приходит с первой битой записью. Обработчик бросил исключение — и многие ждут, что группа встанет. На деле Spring Kafka десять раз подряд повторит ту же запись, напишет строку в лог и поедет дальше: заказ потерян, и узнать об этом можно только из логов.
Дальше — то, что отделяет стенд от production: битая запись, дубли, отставание, смена схемы у семи чужих потребителей и закрытый кластер. Топики, партиции и группы считаются знакомыми по основам Kafka.
Без DLQ Spring Kafka делает по битой записи десять попыток подряд, пишет ошибку в лог и идёт дальше — запись потеряна, и заметно это только в логах. С @RetryableTopic она уходит по retry-топикам с паузами 1, 2 и 4 секунды и через семь секунд ложится в orders-dlt, а основной поток читает партицию дальше.
Подключение: KafkaTemplate, @KafkaListener и момент коммита
Библиотека spring-kafka оборачивает стандартный Java-клиент: KafkaTemplate отправляет, метод с @KafkaListener принимает. Под аннотацией работает контейнер из ConcurrentKafkaListenerContainerFactory: при concurrency = 3 он поднимает три потока, каждый — отдельный потребитель группы со своими партициями. Потоков больше, чем партиций, заводить бессмысленно: лишние получат пустой набор.
@Component
@RequiredArgsConstructor
public class OrderEventPublisher {
private final KafkaTemplate<String, OrderEvent> kafka;
public void publish(OrderEvent event) {
kafka.send("orders", event.orderId().toString(), event);
}
}
@Component
public class OrderEventListener {
@KafkaListener(topics = "orders", groupId = "billing-service")
public void handle(ConsumerRecord<String, OrderEvent> record) {
var event = record.value();
}
}
Главное решение контейнера — когда закоммитить оффсет, то есть сообщить брокеру, что запись обработана. Это AckMode. По умолчанию BATCH: контейнер забирает пачку (до 500 записей), прогоняет через метод и коммитит один раз в конце; упал на трёхсотой — после перезапуска придут все пятьсот. RECORD коммитит после каждой записи: дублей меньше, но каждый коммит — сетевой обмен с координатором группы. MANUAL и MANUAL_IMMEDIATE отдают решение приложению: метод получает Acknowledgment и вызывает acknowledge(), когда запись в базу прошла. MANUAL_IMMEDIATE отправляет коммит сразу, MANUAL копит подтверждения до конца пачки — для метода с @Transactional это решающее: транзакция базы фиксируется на выходе из метода, и MANUAL закоммитит оффсет после неё, а MANUAL_IMMEDIATE — до.
@Bean
ConcurrentKafkaListenerContainerFactory<String, OrderEvent> kafkaListenerContainerFactory(
ConsumerFactory<String, OrderEvent> consumerFactory) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, OrderEvent>();
factory.setConsumerFactory(consumerFactory);
factory.setConcurrency(3);
factory.getContainerProperties().setAckMode(AckMode.MANUAL);
return factory;
}
Повторы и DLQ: куда девать битую запись
Исключение из метода слушателя ловит DefaultErrorHandler — он есть у контейнера всегда, даже необъявленный. Его настройка по умолчанию — FixedBackOff(0, 9): девять повторов без пауз, потом строка ERROR в логе, и оффсет уходит дальше — первый такт схемы выше. Это хуже остановки: остановку видно по отставанию через минуту, потерю находят через неделю по жалобе.
Решение — Dead Letter Queue: после исчерпания попыток запись перекладывается в отдельный топик, оффсет коммитится, основной поток идёт дальше. Вместе с записью уезжают заголовки — текст исключения, стек, исходные топик, партиция и оффсет, — по ним её потом и разбирают.
DefaultErrorHandler с публикацией в DLT
@Bean
DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
var recoverer = new DeadLetterPublishingRecoverer(template,
(record, ex) -> new TopicPartition("orders.DLT", record.partition()));
return new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3));
}
FixedBackOff(1000L, 3) — три повтора с паузой в секунду, четыре попытки всего: второй аргумент считает повторы, а не попытки. Пауза здесь блокирующая: пока обработчик ждёт секунду, партиция стоит вместе со всем, что лежит за битой записью. Три секунды терпимо, минуты — нет.
@RetryableTopic: повторы без остановки партиции
@RetryableTopic(
attempts = "4",
backoff = @Backoff(delayExpression = "${kafka.retry.delay:1000}", multiplier = 2.0),
topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
autoCreateTopics = "false"
)
@KafkaListener(topics = "orders", groupId = "billing-service")
public void handle(OrderEvent event) {
}
@DltHandler
public void handleDlt(OrderEvent event, @Header("kafka_dlt-exception-message") String reason) {
log.error("Не удалось обработать: order={} причина={}", event.orderId(), reason);
}
Здесь повторы не держат партицию: неудачная запись уезжает в orders-retry-0, через секунду — в orders-retry-1, через две — в orders-retry-2, и через семь секунд суммарно ложится в orders-dlt; оффсет в orders закоммичен сразу — второй такт схемы. Без SUFFIX_WITH_INDEX_VALUE топики назывались бы по величине задержки: orders-retry-1000, -2000, -4000. Цена неблокирующих повторов — порядок: пока запись отстаивается в retry-топике, записи того же ключа из основного топика уже обработаны. Где порядок внутри ключа важнее задержки (смена статусов заказа), берут блокирующий DefaultErrorHandler с короткой паузой.
Метод с @DltHandler обязан лежать в том же классе, что и слушатель. Вынесете в соседний бин — Spring подставит заглушку, которая лишь пишет строку в лог, и разбирать orders-dlt будет некому. autoCreateTopics = "false" означает, что retry- и dlt-топики заводят заранее инструментом управления кластером, с тем же числом партиций, что у основного, и с более долгим сроком хранения: записи в DLT ждут человека.
Запись, которую не разобрать
Есть отказ, который ни один обработчик ошибок не увидит: в топик попала запись, которую не может разобрать десериализатор — обрезанный JSON, чужой формат. Падение случается при чтении из брокера, до вызова метода; контейнер читает ту же позицию снова и снова — единственный случай, когда потребитель по-настоящему встаёт. Лечится обёрткой: ErrorHandlingDeserializer ловит исключение настоящего десериализатора и отдаёт запись дальше как битую — её подхватывает обычный обработчик и отправляет в DLT без повторов.
spring:
kafka:
consumer:
key-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
properties:
spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer
Что повторять, а что нет
Повтор имеет смысл только для ошибки, которая пройдёт сама: база не ответила, платёжный шлюз вернул 503. Для неё берут растущую паузу; общий срок — под сценарий: заказу секунды, ночному отчёту минуты. Ошибка, которая повторится при любой попытке — невалидное поле, нарушенное бизнес-правило, — от повторов только множит логи и задерживает соседние записи. Такие исключения объявляют неповторяемыми: errorHandler.addNotRetryableExceptions(IllegalArgumentException.class) у DefaultErrorHandler, exclude = IllegalArgumentException.class у @RetryableTopic — и запись уходит в DLT с первой попытки. Ошибки десериализации и конвертации Spring считает неповторяемыми сам.
Обратный путь из DLT — решение человека: автоматический возврат даёт петлю, где запись с ошибкой в обработчике снова падает и снова ложится в DLT. Поэтому возвращает записи админский эндпоинт с обнулённым счётчиком попыток, а на рост DLT стоит алерт.
Как проверить, что путь до DLT работает
Обработка ошибок — код, который неделями не выполняется, а потом обязан сработать с первого раза. Проверяют его на настоящем брокере: встроенный брокер из spring-kafka-test живёт своей версией, честнее контейнер с тем же образом, что в production. Тест отправляет в orders событие, на котором обработчик бросит неповторяемое исключение, и ждёт его в orders-dlt; паузы укорачивают свойством — ради этого задержка вынесена в ${kafka.retry.delay}.
@SpringBootTest(properties = "kafka.retry.delay=10")
@Testcontainers
class OrderDltTest {
@Container
@ServiceConnection
static ConfluentKafkaContainer kafka = new ConfluentKafkaContainer("confluentinc/cp-kafka:7.8.0");
@Autowired KafkaTemplate<String, OrderEvent> template;
@Test
void invalidOrderLandsInDlt() {
template.send("orders", "order-1", new OrderEvent("order-1", new BigDecimal("-1")));
var props = Map.<String, Object>of(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka.getBootstrapServers(),
ConsumerConfig.GROUP_ID_CONFIG, "dlt-probe",
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
try (var probe = new KafkaConsumer<>(props, new StringDeserializer(), new StringDeserializer())) {
probe.subscribe(List.of("orders-dlt"));
var record = KafkaTestUtils.getSingleRecord(probe, "orders-dlt", Duration.ofSeconds(10));
assertThat(record.headers().lastHeader("kafka_dlt-exception-message")).isNotNull();
}
}
}
DLT читает отдельный потребитель со строковым десериализатором: слушатель приложения с JsonDeserializer споткнулся бы на битой записи ещё раз. Брокер в контейнере создаёт топики сам, так что autoCreateTopics = "false" тесту не мешает.
Дубли: идемпотентный потребитель
Kafka обещает доставку хотя бы раз, и «хотя бы» — не оговорка. Оффсет коммитится после обработки: упали между ними — запись придёт снова. При каждом деплое группа перераспределяет партиции, и незакоммиченную пачку новый владелец читает заново. Запись, возвращённая из DLT, приходит второй раз по определению. Значит, обработчик «оплата получена» рано или поздно увидит одну оплату дважды и, если просто начисляет, начислит дважды.
Защита — таблица обработанных событий с уникальным ключом, но порядок действий важнее таблицы. Напрашивается «проверить, не обработано ли, и если нет — обработать»; с двумя экземплярами это не работает: оба видят «отметки нет», оба начисляют. Правильный порядок — сначала вставить отметку, потом работать, в одной транзакции базы: одновременную вставку пропустит только один, второй споткнётся об уникальный индекс, а ON CONFLICT DO NOTHING превратит отказ в ноль строк.
@KafkaListener(topics = "orders", groupId = "billing-service")
@Transactional
public void handle(OrderEvent event, Acknowledgment ack) {
if (processedEvents.insertIfAbsent(event.eventId(), "billing-service") == 0) {
ack.acknowledge();
return;
}
billing.charge(event.orderId(), event.amount());
ack.acknowledge();
}
Ключ — eventId, который продьюсер выдаёт событию в своей транзакции (у outbox это идентификатор строки), а не оффсет: он свой у каждой группы и меняется при переразбиении топика. Если в одной базе живут слушатели нескольких групп, ключ составной — (event_id, consumer_group), иначе отметка первого спрячет событие от второго. Таблицу чистят по processed_at, но не раньше, чем перестанут приходить повторы: возврат из DLT случается и через неделю, поэтому хранят дни. И транзакция базы закрывает только базу: при вызове платёжного шлюза по HTTP повтор пройдёт мимо таблицы — шлюзу передают тот же eventId заголовком Idempotency-Key. Разбор гонки двух экземпляров — в распределённых паттернах.
Транзакции Kafka и read_committed
Конвейер «прочитал из orders, обогатил, записал в orders-enriched» ломается иначе. Потребитель отправил результат и упал до коммита оффсета — после перезапуска отправит второй раз, и потребители ниже по течению должны гасить дубль сами. Kafka закрывает это окно транзакцией: отправка в выходной топик и коммит оффсета входного фиксируются вместе или откатываются вместе.
spring:
kafka:
producer:
transaction-id-prefix: billing-tx-
consumer:
isolation-level: read_committed
С префиксом Spring Boot создаёт KafkaTransactionManager, и контейнер оборачивает в транзакцию каждую пачку: всё отправленное через KafkaTemplate и оффсеты обработанных записей уходят одним коммитом, исключение из метода откатывает и то и другое. Продьюсер с тем же transactional.id при повторном запуске выбивает старого — так отсекаются «зомби», дописывающие транзакцию после перезапуска. Читателям выходного топика нужен read_committed, иначе они видят записи незавершённых и откаченных транзакций.
Границы гарантии — стены Kafka: запись в базу или HTTP-вызов внутри метода в транзакцию не входят, для них остаётся идемпотентный потребитель. Цена: задержка выходного топика растёт до размера пачки; логу __transaction_state по умолчанию нужны три брокера (transaction.state.log.replication.factor=3, min.insync.replicas=2), и на стенде с одним брокером приложение не стартует, пока их не понизить. Самое неприятное: пока одна транзакция открыта, все read_committed-читатели партиции стоят на её границе — до transaction.timeout.ms, по умолчанию минуты. Поэтому транзакции берут для «топик → топик», а для «топик → база» — идемпотентность.
Отставание: consumer lag, ребаланс и обратное давление
Событие записано, а действия по нему нет: письмо не ушло, склад не зарезервировал. Продьюсер здоров, потребитель не падает — читает медленнее, чем пишут. Расстояние между концом партиции и позицией группы называется lag, и у него две природы: 80 записей после всплеска — норма, рост при стоящей позиции — авария.
Lag — расстояние от позиции группы до конца партиции, и само по себе число ничего не значит: 80 записей после всплеска исчезнут через секунду. Тревога — когда конец журнала уходит, а позиция стоит: потребитель либо упал, либо раз за разом повторяет одну и ту же запись.
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group billing-service --describe
# TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# orders 0 124530 124530 0
# orders 1 124450 124530 80 ← отстаёт
Причин четыре. Медленный обработчик — синхронный запрос к базе или во внешний сервис на каждую запись; лечится пачками и паузой, о них ниже. Нехватка партиций — потребителей не добавить больше, чем партиций. Истёкший max.poll.interval.ms — пачка обрабатывалась дольше пяти минут, брокер счёл потребителя мёртвым и запустил ребаланс; со стороны это «обработка идёт, а lag растёт»: группа по кругу перечитывает одну пачку. Паузы JVM на сборке мусора — потребитель молчал тридцать секунд, lag вырос на тысячи разом.
Потолок пропускной способности задают партиции
Если один потребитель обрабатывает 200 записей в секунду, а партиций восемь, потолок группы — 1 600 в секунду: девятый потребитель получит пустой набор. На пике в 5 000 в секунду lag растёт на 3 400 в секунду, за час — на 12 миллионов. Партиции можно добавить, но ключи лягут иначе: события одного заказа из третьей партиции поедут в одиннадцатую, и на время перехода порядок внутри заказа теряется. Поэтому партиции считают при создании топика с двукратным запасом к пику (48 дают 9 600 в секунду): лишние дёшевы, пересчёт ключей дорог. И у событий должен быть срок годности: продьюсер кладёт expiresAt, потребитель отбрасывает протухшее с записью в метрику вместо письма о заказе, доставленном два часа назад.
Ребаланс: почему деплой создаёт дубли и как его укоротить
Когда потребитель входит в группу или выходит, брокер переразбирает партиции. В классическом протоколе у всех отбираются все партиции и раздаются заново, а незакоммиченная пачка каждого достаётся новому владельцу повторно — так каждый деплой рождает дубли и секунды простоя. Скорость реакции задают две настройки. session.timeout.ms (45 секунд с Kafka 3.0) — сколько молчания без heartbeat, чтобы потребителя сочли мёртвым; max.poll.interval.ms (5 минут) — сколько без poll(), чтобы сочли зависшим. Первую снижают, чтобы быстрее заметить упавший под; вторую поднимают, когда обработка честно долгая, или укорачивают пачку через max.poll.records.
Остановку целиком убирает кооперативный ребаланс: partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor — отбираются только партиции, которые правда переезжают, остальные читают дальше. От перезапуска пода спасает статическое членство group.instance.id=billing-1: вернувшийся под тем же именем в пределах session.timeout.ms получает свои партиции без ребаланса; при concurrency > 1 Spring дописывает к имени номер потока. В Kafka 4.0 есть новый протокол групп (group.protocol=consumer): раскладку считает брокер, ребаланс перестаёт останавливать группу — нужны клиент и брокер версии 4.
Пачки и пауза: когда обработка не успевает
Пятьсот записей — пятьсот INSERT, а мог быть один. Пакетный слушатель принимает всю пачку:
@KafkaListener(topics = "orders", groupId = "billing-service", batch = "true")
public void handle(List<OrderEvent> events) {
billing.chargeAll(events);
}
Тонкость с ошибками: обычное исключение повторит все пятьсот записей. Чтобы отправить в DLT только виноватую, бросают BatchListenerFailedException с её индексом — обработчик закоммитит всё до неё и повторит с неё.
Когда не успевает не обработчик, а система за ним — база под нагрузкой, шлюз отвечает ошибками, — повторы её добивают. Правильная реакция — перестать читать: registry.getListenerContainer("orders").pause() через KafkaListenerEndpointRegistry. Контейнер на паузе продолжает вызывать poll(), членство в группе живо и ребаланса нет, но записи в метод не идут; resume() — когда предохранитель (Circuit Breaker) перед шлюзом закрылся. Для одной записи то же делает ack.nack(Duration.ofSeconds(5)).
Наблюдаемость: что смотреть кроме lag
У lag два имени из разных мест. Клиент Kafka отдаёт отставание через JMX, и в Spring через Micrometer оно приезжает как kafka_consumer_fetch_manager_records_lag — по каждой партиции, глазами приложения. kafka_consumergroup_lag из готовых дашбордов принадлежит стороннему экспортёру (kafka_exporter или Burrow): его ставят отдельно, зато он видит отставание и когда приложение лежит. Алерт один — lag больше порога и растёт N минут подряд; порог зависит от топика: для платежей сотни, для аналитики сотни тысяч.
Записи не говорят, насколько вы опоздали: 40 000 — это десять секунд при 4 000 в секунду и два часа при пяти. Поэтому рядом держат возраст обрабатываемой записи, now - record.timestamp(), и алертят на него как на отставание во времени. Третья метрика — таймер spring_kafka_listener_seconds: его 99-й перцентиль растёт раньше lag. Четвёртая — DLT: любая запись в DLT платёжного топика — инцидент, а не статистика; с экспортёром это increase(kafka_topic_partition_current_offset{topic=~".*-dlt"}[30m]) > 0.
Трассировка идёт через заголовки — пары «ключ, байты» рядом с телом записи, место для технического, а не бизнес-данных. Spring Boot 3 с Micrometer Tracing сам пишет traceparent при отправке и продолжает трассу в слушателе, если включить spring.kafka.template.observation-enabled и spring.kafka.listener.observation-enabled. Свои заголовки читают через @Header("X-Source-Service") String source.
Schema Registry: как не сломать соседей
Команда заказов переименовала в событии поле amount в total и выкатилась. Семь потребителей не упали: JsonDeserializer в Spring не ругается на незнакомые поля, amount стал null, биллинг начал выставлять счета на ноль — заметили по выручке через два дня. С Avro без реестра ломается громче: у потребителя нет схемы, по которой писали, и каждая запись падает при чтении — DLT заполняется тысячами записей за минуту. Корень один: договор о формате нигде не записан, и проверить его до выкладки нечем.
Schema Registry — сервис, где схема живёт отдельно от сообщений, а каждая новая версия проходит проверку на совместимость до того, как первое сообщение с ней уйдёт в топик.
Схема не едет в каждом сообщении: продьюсер один раз регистрирует её и получает номер, а в записи остаются пять служебных байт и голые данные. Потребитель по номеру забирает схему из реестра, кэширует и декодирует байты. Несовместимую версию реестр отбивает на шаге 1 — до того, как первое сообщение с ней уйдёт в топик.
Продьюсер при первой отправке регистрирует схему и получает номер; в каждую запись он пишет пять служебных байт — ноль как признак формата и четырёхбайтный schema_id — и дальше данные без имён полей. Потребитель читает номер, один раз забирает схему из реестра, кэширует и декодирует.
Стандарт экосистемы — Avro: компактный бинарный формат с формальными правилами эволюции; классы по .avsc генерирует avro-maven-plugin при сборке. Protobuf берут команды, у которых он уже есть ради gRPC. JSON Schema читается глазами и в разы объёмнее — для отладки и малых потоков.
Совместимость реестр проверяет по режиму топика. BACKWARD (по умолчанию): потребитель с новой схемой обязан читать старые записи — можно удалять поля и добавлять необязательные с default, нельзя добавить обязательное. Переименование — удаление плюс добавление обязательного, и реестр отвечает 409 Conflict: продьюсер падает на первой отправке, а не потребители через два дня. FORWARD — зеркало: потребитель со старой схемой обязан читать новые записи. FULL — оба условия, NONE — без проверок. Режим выбирают по тому, кого обновляют первым: BACKWARD требует обновить потребителей раньше продьюсера и годится, когда потребители ваши или историю топика перечитывают новой схемой. Семь чужих сервисов, выкатываемых вразнобой, — FORWARD или FULL: старый потребитель обязан пережить новое событие. И правило production: auto.register.schemas=false у продьюсера, схемы регистрирует конвейер сборки — несовместимая версия не доезжает до старта сервиса.
Тюнинг: продьюсер и потребитель
Из коробки Kafka настроена на минимальную задержку: продьюсер отправляет запись сразу, потребитель просит данные, даже если у брокера один байт. При десятках тысяч записей в секунду это тысячи мелких запросов; пропускную способность покупают задержкой в миллисекунды.
Продьюсер
| Параметр | Значение по умолчанию | Когда менять |
|---|---|---|
batch.size | 16 KB | Поднять до 64–256 KB при большом объёме записи: продьюсер копит пачку на партицию и отправляет её целиком. |
linger.ms | 0 (5 мс с Kafka 4.0) | Поднять до 5–20 мс: продьюсер ждёт столько перед отправкой неполной пачки. При 1 000 записей в секунду 10 мс собирают по десять записей в запрос вместо одной. |
compression.type | none | Включить zstd (Kafka 2.1+) или lz4: на JSON экономит в 3–5 раз трафик и диск. Сжимает продьюсер, брокер хранит как есть, распаковывает потребитель — процессор тратится на концах, не на брокере. |
Потребитель
| Параметр | Значение по умолчанию | Когда менять |
|---|---|---|
fetch.min.bytes | 1 байт | Поднять до 50–500 KB: брокер отвечает, когда накопил столько или истёк fetch.max.wait.ms (500 мс) — меньше пустых запросов. |
max.poll.records | 500 | Снизить до 50–100, если запись обрабатывается долго: короче пачка — меньше повтор при сбое и дальше от max.poll.interval.ms. |
max.poll.interval.ms | 5 минут | Поднять, если обработка пачки честно занимает столько; иначе ребаланс по кругу. |
Потолок размера записи
Запись больше мегабайта Kafka не примет. Мегабайт — это JSON примерно на десять тысяч строк заказа или одна отсканированная страница. Лимит стоит с трёх сторон: max.request.size у продьюсера (1 048 576), message.max.bytes у брокера (1 048 588, на топик — max.message.bytes) и max.partition.fetch.bytes у потребителя. Продьюсер меряет запись до сжатия и отвечает RecordTooLargeException, брокер — после, поэтому сжатие помогает пройти брокера, но не продьюсера. Поднимают все три: продьюсер и брокер — чтобы запись прошла, потребитель — чтобы не читать по одной, ведь пачка больше его лимита доедет, но одна за запрос. Файлы, фото и PDF в Kafka не кладут: в объектное хранилище, а в событие — ключ.
Безопасность: кто ты и что тебе можно
По умолчанию Kafka слушает порт 9092 без шифрования и проверки клиента: любой, кто дотянулся до сети, читает orders консольным потребителем и пишет в payments. Защита отвечает на три вопроса, и каждый слой без соседних бесполезен.
Первый — не подслушают ли и тот ли это брокер. Это TLS: трафик шифруется, а клиент проверяет, что имя в сертификате брокера совпадает с адресом подключения (ssl.endpoint.identification.algorithm=https, по умолчанию с Kafka 2.0); без проверки имени подменить брокер по дороге можно и с шифрованием. Второй — кто клиент. Это SASL: PLAIN передаёт логин и пароль внутри TLS, брокер хранит пароль открытым — для стенда; SCRAM-SHA-256 и SCRAM-SHA-512 — обмен вызов-ответ с солёным хешем, пароль по сети не идёт и на брокере не лежит; OAUTHBEARER — токены провайдера идентификации, обновляются без перезапуска. Третий — что клиенту можно. Это ACL. Аутентификация без авторизации показывает, кто читает payments, но не мешает: любой сервис, прошедший SASL, видит все топики. С включённым авторизатором запрещено всё, что не разрешено явно (allow.everyone.if.no.acl.found=false), и у каждого сервиса свой пользователь с правами ровно на свои топики:
kafka-acls.sh --bootstrap-server kafka:9093 --add \
--allow-principal User:billing-service \
--operation Read --topic orders --group billing-service
Право на группу забывают чаще всего: без него потребитель падает на старте с GroupAuthorizationException, хотя топик читать ему можно.
spring.kafka.properties.security.protocol=SASL_SSL
spring.kafka.properties.sasl.mechanism=SCRAM-SHA-256
spring.kafka.properties.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
username="billing-service" password="${KAFKA_PASSWORD}";
spring.kafka.properties.ssl.truststore.location=/etc/kafka/truststore.jks
spring.kafka.properties.ssl.truststore.password=${TRUSTSTORE_PASSWORD}
Пароль вынесен в переменную окружения, а truststore «просто лежит по пути». Внутри не сертификаты брокеров, а сертификат удостоверяющего центра, который их подписал: клиент доверяет центру, а значит, любому брокеру с его подписью. Файл публикует команда платформы; на машину он попадает секретом Kubernetes, смонтированным по /etc/kafka/truststore.jks, или выгружается из хранилища секретов при старте; если корпоративный центр уже лежит в cacerts JVM, отдельный truststore не нужен. Сертификаты брокеров cert-manager перевыпускает каждые 60–90 дней, и клиенты этого не замечают — подпись та же. Центр меняется раз в годы: новый сертификат кладут рядом со старым заранее, выкатывают клиентов, потом убирают старый. Truststore клиент читает при создании фабрики продьюсера или потребителя, так что подменить файл мало — нужен перезапуск.
KRaft: где живут метаданные
Метаданные кластера — какие есть топики, кто лидер каждой партиции, какие действуют ACL — до Kafka 3.3 хранил отдельный кластер ZooKeeper: вторая система со своими портами, мониторингом и режимами отказа.
Метаданные — какие есть топики и кто лидер каждой партиции — раньше лежали в отдельном ZooKeeper, и его отказ был отдельным режимом отказа Kafka. В KRaft они хранятся таким же реплицированным журналом, как данные, а кворум контроллеров держат сами брокеры.
KRaft переносит метаданные внутрь Kafka: несколько брокеров работают контроллерами, держат кворум по протоколу Raft и хранят метаданные журналом __cluster_metadata, реплицированным как любой топик. В маленьком кластере брокер и контроллер — один процесс (process.roles=broker,controller), в большом контроллеры выделяют. Появился KRaft в 2.8 как ранний доступ, готовым к production объявлен в 3.3, в Kafka 4.x поддержка ZooKeeper убрана: кластер на ZooKeeper сначала мигрируют на KRaft в 3.x, потом обновляют. Для кода приложения ничего не меняется — bootstrap-servers и клиент те же.
Глубже: что фреймворк делает сам: десять попыток и пропуск записирасширенное
Раздел про DLT показывает, как настроить обработку ошибок. Важно знать и то, что происходит, если её не настраивать: умолчания Spring Kafka и брокеров за последние версии изменились, и статьи со «бесконечным циклом» описывают старое поведение.
У контейнера слушателя Spring Kafka DefaultErrorHandler есть всегда, и без настройки он повторяет запись десять раз подряд без паузы (FixedBackOff(0, 9)), после чего пишет ошибку в лог и переходит к следующей записи. Цикла нет, но нет и DLT: запись теряется, а лог это единственный след. Для пакетного слушателя повторяется вся пачка, и одна битая запись десять раз прогоняет через обработчик соседей, что для неидемпотентного обработчика опасно. Первое, что меняют, это восстановитель (DeadLetterPublishingRecoverer) и пауза между попытками; список исключений, которые не повторяют (addNotRetryableExceptions), второе.
У RabbitMQ зеркальная история: надёжные очереди с версии 4.0 сами ограничивают доставки двадцатью (x-delivery-limit), после чего сообщение уходит в DLX или удаляется, если DLX нет, а классические очереди по-прежнему крутят его вечно, о чём статья про RabbitMQ в production.
Итог для проверки в новом проекте: где окажется запись после последней попытки, в DLT, в DLX или в логе; повторяется ли то, что повторять бессмысленно (ошибка разбора и валидации); и есть ли тест, который это доказывает, как в разделе выше. Ответ «фреймворк разберётся» верен только в том смысле, что он не зациклится; куда денется запись, решаете вы.
Коротко
- Битая запись по умолчанию не останавливает группу, а теряется после десяти попыток: DLT через
@RetryableTopicобязателен,ErrorHandlingDeserializer— для неразбираемых записей, путь до DLT проверяет тест на контейнере. - Повторяют только то, что пройдёт само; невалидное — в DLT с первой попытки. Возврат из DLT — кнопка для человека, не цикл.
- «Хотя бы раз» — дубли при каждом деплое: отметка по
eventIdвставляется до работы и в той же транзакции; гонку двух экземпляров решает уникальный индекс. - Транзакции Kafka и
read_committedзакрывают дубли только внутри Kafka; для базы и HTTP остаётся идемпотентность. - Партиции — потолок группы, их считают при создании топика с запасом; рост lag при стоящей позиции — авария. Ребаланс укорачивают
CooperativeStickyAssignorиgroup.instance.id, перегрузку снимают пачками иpause(). Кроме lag смотрят возраст записи, время в слушателе и рост DLT. - Schema Registry ловит несовместимую схему при регистрации, а не у потребителей; чужие потребители — FORWARD или FULL.
- Лимит записи 1 МБ стоит с трёх сторон; файлы в Kafka не кладут.
- TLS — «тот ли брокер», SASL — «кто клиент», ACL — «что ему можно»; без ACL любой аутентифицированный сервис читает всё. В truststore — центр, а не брокеры.
- KRaft: метаданные в самом кластере, с 4.x ZooKeeper нет.
- Без настройки Spring Kafka повторяет запись десять раз без паузы и пропускает её с записью в лог, а не в DLT; у надёжных очередей RabbitMQ предел двадцать доставок; куда попадёт запись после последней попытки, решаете вы.
Что пощупать
Момент коммита, повтор доставки и идемпотентный потребитель в практикуме remodov/marketplace-system реализованы дважды и одинаково: сервис заказов принимает PaymentCompleted, сервис уведомлений — события заказа, и оба сначала пишут идентификатор события в журнал обработанных, а уже потом делают работу.
Код: adapter-in-kafka, JooqProcessedEventsRepository, потребитель уведомлений: OrderEventConsumer.
Что почитать дальше
- Основы Kafka — партиции, ключи, гарантии доставки, outbox: фундамент этой статьи.
- Распределённые паттерны — outbox и Idempotent Consumer подробно, с гонкой двух экземпляров.
- Kafka Streams —
exactly_once_v2: транзакции Kafka в полную силу для конвейеров «топик → топик». - Стандарт Kafka — те же правила чек-листом для ревью: retry-топики, идемпотентность, наблюдаемость, безопасность.