Привет, Хабр! Меня зовут Максим Шуматбаев, я инженер технической поддержки в Arenadata. Эта статья — попытка собрать в одном месте внутренние механизмы Apache Kafka, которые задействованы при работе клиентских приложений. Цель простая: дать понимание, с которым типовые инциденты перестают выглядеть загадкой и разбираются за минуты, а не за дни.

Типичный инцидент. Сервис месяцами работал без каких-либо проблем, и вот мониторинг показывает растущий лаг. В логах куча ошибок, клиентское приложение будто потеряло кластер, но тот же мониторинг показывает, что брокеры живы. Сопровождение сервиса говорит, что проблема в кластере, администраторы Kafka парируют и винят во всем приложение. Спустя время все приходит в норму само по себе. Проблема решена, но как именно никто не понимает, потому что мало кто представляет, что вообще происходит между вызовом producer.send() и моментом, когда сообщение прочитает другой сервис.

А происходит там много. Брокеры Kafka сами по себе довольно простые — они хранят логи и отвечают на клиентские запросы. Большая часть логики живет именно в клиентских библиотеках - клиент сам выбирает, какому брокеру слать запись, сам решает, когда обновить карту кластера, сам собирает батчи, сам ретраит при сбоях, сам держит порядок и фиксирует прогресс чтения. Понимать поведение Kafka во много означает понимать поведение клиента, а не брокера.

Ниже разобран весь путь сообщения, а именно: как клиент находит нужный брокер, что реально происходит при записи, как кластер защищается от потерь и дублей, как работают транзакции, как потребитель читает и коммитит offset и какие эксплуатационные ограничения умеют тихо все замедлить. После этого «растущий лаг без причины», как правило, перестает быть загадкой, почти всегда за ним стоит один из описанных механизмов. Материал получился объемным, поэтому разбит на две части: в первой — про метаданные и про запись, во второй — про чтение, consumer group и эксплуатационные ограничения.

Все сказанное относится к ванильной Apache Kafka и ее штатным клиентам. Конфиги и команды работают на любом совместимом дистрибутиве.

1. Вспомним основы

Модель хранения — это фундамент, на него опирается все остальное. Пройдемся быстро.

Broker, Topic, Partition

Kafka хранит данные на серверах-брокерах и раскладывает их по именованным темам (topic). Топик режется на разделы (partition), которые распределены по брокерам и дискам ради масштабирования и отказоустойчивости. Партиция —это упорядоченный неизменяемый лог — записи только дописываются в конец, не меняются и не удаляются. Каждой записи Kafka присваивает порядковый номер (offset). Отсюда и производительность — последовательные операции записи и чтения самые дешевые операции для диска.

Хранение сообщений в партициях Kafka. На рисунке показан топик, разбитый на несколько партиций, где каждое сообщение внутри раздела получает offset. Новые данные всегда добавляются в конец партиции
Хранение сообщений в партициях Kafka. На рисунке показан топик, разбитый на несколько партиций, где каждое сообщение внутри раздела получает offset. Новые данные всегда добавляются в конец партиции

Leader, Follower, ISR

Каждая партиция реплицируется на несколько узлов. У нее один лидер (leader) и ноль или больше фолловеров (follower), копий на других брокерах. Клиенты ходят только к лидеру: продюсеры пишут в него, потребители из него читают. Фолловеры подтягивают записи с лидера и держат себя в актуальном состоянии.

Те фолловеры, что успевают реплицировать данные, считаются синхронными репликами и вместе с лидером образуют ISR (In-Sync Replicas). Сильно отстал, тебя временно выкинули из ISR. Догнал — вернули. Упал лидер, контроллер выбирает нового из того, что осталось в ISR.

Репликация данных в Kafka. На рисунке показан кластер из трех брокеров с фактором репликации 3. Продюсер отправляет сообщение лидеру партиции (красный узел), после чего лидер асинхронно реплицирует запись на фолловеры (синие узлы)
Репликация данных в Kafka. На рисунке показан кластер из трех брокеров с фактором репликации 3. Продюсер отправляет сообщение лидеру партиции (красный узел), после чего лидер асинхронно реплицирует запись на фолловеры (синие узлы)

Запомните связку «лидер — ISR-подтверждения записи», к ней мы вернемся в главе про надежность. Именно от ISR зависит, переживет ли запись падение брокера.

2. Метаданные — наше все

Как клиент понимает, куда идти

Первая проблема — клиент не знает заранее, какой брокер сейчас лидер нужной партиции. И не может знать, потому что все меняется: брокеры падают, лидеры переизбираются, узлы добавляются, поэтому жестко прописать лидеров в конфиг невозможно. За подобные запросы отвечает Metadata API.

Метаданные — это текущее состояние кластера: все брокеры с адресами, все топики с партициями и, главное, кто сейчас лидер каждой партиции. Брокеры эту информацию только хранят и выдают, а к кому идти решает сам клиент.

ZooKeeper больше не нужен. Раньше метаданные и выборы лидеров координировал ZooKeeper. С Kafka 4.0 его убрали совсем, кластер работает в режиме KRaft, где координацию берет на себя внутренний кворум контроллеров на самих брокерах. Для клиента ничего не меняется, он по‑прежнему просит у брокеров карту кластера, просто источник этой карты стал другим.

Чтобы получить такую актуальную карту, хватит одного брокера. В конфиге задается bootstrap.servers, список из одного или нескольких адресов для первого знакомства с кластером. Клиент идет по списку и пробует установить соединение — не ответил первый, берет следующий. На практике одного живого брокера достаточно, чтобы узнать про остальные, однако несколько адресов держат на случай недоступности брокеров. Подключились, сразу шлем MetadataRequest и получаем полную картину.

Обмен данными при подключении клиента к Kafka: клиент устанавливает соединение с брокером, проходит аутентификацию, затем отправляет MetadataRequest. Брокер читает состояние кластера из контроллера/ZooKeeper и отвечает MetadataResponse, содержащий ID контроллера, список брокеров (NodeId и адреса), а также список топиков с информацией о разделах и указанием лидера каждого раздела. Клиент заполняет свой внутренний кэш полученными метаданными
Обмен данными при подключении клиента к Kafka: клиент устанавливает соединение с брокером, проходит аутентификацию, затем отправляет MetadataRequest. Брокер читает состояние кластера из контроллера/ZooKeeper и отвечает MetadataResponse, содержащий ID контроллера, список брокеров (NodeId и адреса), а также список топиков с информацией о разделах и указанием лидера каждого раздела. Клиент заполняет свой внутренний кэш полученными метаданными

В MetadataResponse приходит снимок топологии. Самое важное в нем — это лидер каждой партиции. Если запрос по ошибке прилетит не к лидеру, брокер ответит NotLeaderForPartition, это означает, что клиентский запрос не по адресу.

Дальше клиент кэширует карту в памяти и маршрутизирует запросы сам. Это ключевая архитектурная идея, в Kafka нет центрального прокси, через который течет весь трафик. Продюсер знает, что запись для topic-a:partition-0 идет на broker-1, и открывает соединение прямо туда. Потребитель сам шлет FetchRequest лидеру своей партиции. Нагрузка размазана по всем лидерам, брокеры не занимаются лишним проксированием, что повышает производительность всех операций на запись-чтение, а также не создает какой-либо единой точки отказа.

Использование кэша метаданных на стороне клиента для маршрутизации запросов. Продюсер и консюмер обращаются к локальному кэшу: (1) продюсер запрашивает в кэше лидера для topic-a:partition-1 и получает broker-1, после чего отправляет сообщение напрямую этому брокеру; (3) консюмер запрашивает в кэше лидера для topic-c:partition-2 – получает broker-0 и затем делает FetchRequest напрямую broker-0. И продюсер, и консюмер поддерживают открытые подключения к тем брокерам, которые являются лидерами нужных им разделов.
Использование кэша метаданных на стороне клиента для маршрутизации запросов. Продюсер и консюмер обращаются к локальному кэшу: (1) продюсер запрашивает в кэше лидера для topic-a:partition-1 и получает broker-1, после чего отправляет сообщение напрямую этому брокеру; (3) консюмер запрашивает в кэше лидера для topic-c:partition-2 – получает broker-0 и затем делает FetchRequest напрямую broker-0. И продюсер, и консюмер поддерживают открытые подключения к тем брокерам, которые являются лидерами нужных им разделов.

Получается простой и надежный service discovery: любой брокер выдает полную карту, достаточно знать про один. При масштабировании и отказах клиент узнает про новые брокеры, перераспределение партиций и смену лидеров автоматически.

Когда клиент решает обновить кэш

Что мы имеем —  клиент закэшировал топологию, считает ее верной и какое-то время не трогает. Но в распределенной системе ничто не вечно: на каком-то брокере решили провести работы, из-за чего он недоступен, у партиции сменился лидер, в новом релизе приложения появились новые топики. Кэш устаревает, и его надо обновить (refresh).

Правило простое: доверяй кэшу, пока не получишь ошибку. Клиент перестает верить метаданным, только когда узнает, что они врут. Сигналом служат ошибки от брокеров или внутренний таймер. Основные ошибки-триггеры:

  • NotLeaderForPartition. Брокер больше не лидер нужной партиции. Клиент помечает кэш недействительным и обновляет его.

  • UnknownTopicOrPartition. Брокер не знает такого топика или партиции. Например, продюсер пишет в свежесозданный топик, про который информация еще не разошлась.

  • Сбой соединения. TCP оборвалось, клиент считает брокер недоступным и идет за метаданными к другому узлу.

Кроме реактивного обновления, есть проактивное, периодический refresh на всякий случай с интервалом metadata.max.age.ms. Даже без ошибок клиент сам сходит за свежей топологией. Это позволит увидеть изменения, которые ошибкой не проявляются —  например, расширили число партиций в существующем топике, ошибки не будет, пока клиент не обратится к новым партициям, и только периодический refresh гарантирует, что он про них узнает.

Обновление метаданных: сочетание реактивного и проактивного подхода. В нормальном состоянии клиент считает метаданные актуальными и направляет запросы к leader-партициям. При этом на фоне тикает таймер metadata.max.age.ms, по истечении которого наступает проактивный триггер. Если же брокер отвечает ошибкой, возникает реактивный триггер. Оба пути ведут к запросу обновленных метаданных.
Обновление метаданных: сочетание реактивного и проактивного подхода. В нормальном состоянии клиент считает метаданные актуальными и направляет запросы к leader-партициям. При этом на фоне тикает таймер metadata.max.age.ms, по истечении которого наступает проактивный триггер. Если же брокер отвечает ошибкой, возникает реактивный триггер. Оба пути ведут к запросу обновленных метаданных.

Но во всем важен баланс, в том числе и частоте таких обновлений. Редкие refresh - клиент дольше работает по устаревшей карте и ловит временные NotLeader и UnknownPartition. Частые refresh —  лишняя нагрузка на брокеры и рост задержек. В норме чаще, чем раз в несколько минут, обновляться незачем.

Симптомы проблем маршрутизации

Большая часть непонятных проблем с Kafka, как правило, возникает из-за того, что клиенты не туда или не так шлют свои запросы к Kafka. Все это выглядит как нестабильный кластер, а на самом деле упирается в сеть или неправильную конфигурацию на клиентах. Разберем, что обычно стоит за каждым симптомом.

Если клиент вообще не подключается к кластеру, причина почти всегда в bootstrap.servers или в сети: указанный брокер недоступен, а запасных адресов в списке нет. Частый случай — это когда клиент бесконечно стучится в какой-то неизвестный хост. Здесь почти наверняка виноваты listeners или advertised.listeners на брокере, об этом отдельно ниже.

Дальше идет семейство ошибок, которые обычно означают одно и то же: клиент работает по устаревшей карте кластера. Постоянные NotLeaderForPartition говорят, что клиент стучится не к тому брокеру: либо сеть мешает ему обновить метаданные, либо кластер реально нестабилен и лидеры успевают смениться раньше, чем применится свежий кэш. UnknownTopicOrPartition означает, что брокер не знает запрошенный топик или партицию, например, продюсер пишет в несозданный топик при выключенном автосоздании или указан неверный номер партиции. Для только что созданного топика это нормально при первом обращении, клиент сам сходит за метаданными. LeaderNotAvailable означает, что у партиции временно нет лидера, обычно сразу после создания топика или во время перевыборов.

Общее правило, по которому можно определить настоящие проблемы, простое. Единичные всплески этих ошибок сразу после перезагрузки брокера или изменения числа партиций — это ожидаемо: клиент ловит ошибку в момент перестроения и за секунды восстанавливается. А вот если ошибки не разовые, а постоянные, дело уже не в перестроении, а в сети или конфигурации, и искать нужно именно в этом направлении.

Важный момент, advertised.listeners. Брокер сообщает клиентам свои адреса через метаданные. Если кластер живет в сложной сети (Docker, Kubernetes, NAT, облако), нужно правильно прописать advertised.listeners, адреса, по которым клиент реально сможет достучаться до кластера. Иначе брокер честно отдаст свой внутренний адрес вроде PLAINTEXT://kafka.internal:9092, снаружи недоступный, и клиент будет стучаться в пустоту. Разница такая: listeners — это интерфейсы, на которых брокер слушает, advertised.listeners — это то, что он сообщает клиентам для подключения.

3. Как на самом деле происходит запись

Карта есть, лидеры известны, дальше сама запись. Любое событие, отправляемое в Kafka, проходит несколько стадий. Каждый этап — это отдельный набор механизмов, поэтому имеет свое влияние на задержку и нагрузку.

Путь сообщения: этапы и их смысл

Сериализация. Продюсер кодирует ключ и значение в байты. На данном этапе наиболее важным является выбранный формат кодирования (Avro, Protobuf). Правильный формат позволит экономить и процессорное время, и выходной объем байт. До сериализации опционально отрабатывают интерцепторы (Interceptors), которые позволяют применить дополнительную логику над сообщением, например, настроить дополнительный аудит.

Выбор партиции. Если вместе с сообщением передается ключ, то продюсер хеширует его (по умолчанию murmur2) и берет остаток от деления на число партиций. Все записи с одним ключом попадают в одну партицию, порядок по ключу сохраняется. Если же ключа нет, то работает липкий (sticky) механизм: продюсер закрепляется за случайной партицией и льет сообщения туда, пока не наберется батч или не истечет linger.ms, потом происходит переключением на другую партицию.

Буферизация и батчинг. Выбрав партицию, продюсер кладет сообщение в буфер RecordAccumulator. Под каждую пару топик-партиция своя двухсторонняя очередь из ProducerBatch, все батчи делят общий пул памяти из buffer.memory. Использование батчей позволяет размазать накладные расходы протокола, что позволяет достичь той самой высокой пропускной способности Kafka.

Сжатие. Батч можно сжать перед отправкой, поддерживаются gzip, snappy, lz4, zstd. Сжатие идет на уровне батча, поэтому чем крупнее пачка сообщений, тем лучше коэффициент сжатия. Компрессия экономит сеть и диск ценой CPU, и обычно выигрыш это окупает.

Отправка. Поток Sender забирает готовые батчи и шлет брокерам. Важная деталь: продюсер держит одно соединение на брокер, а не на партицию. Например, если брокер выступает лидером пяти нужных партиций, то продюсер отправляет все пять потоков через один сокет. Это позволяет сэкономить на установке соединений (особенно заметно с SSL).

Один запрос, много партиций. ProduceRequest несет сразу несколько батчей для разных партиций, если у них общий лидер. Продюсер группирует записи по брокерам: один пакет на topic-a со всеми его партициями, другой на topic-b. Это повышает эффективность записи за счет меньшего количества сетевых запросов.

Подтверждение (ack). Записав батч брокер возвращает ProduceResponse. Сколько ждать подтверждения, зависит от acks: от «не жду вообще» до «жду все ISR». В ответе на каждую партицию приходит либо offset, либо ошибка.

Партиции: как выбор партиции влияет на порядок и нагрузку

Партиционирование один из фундаментальных механизмов Kafka. Как продюсер раскидывает записи по партициям, так и балансируется нагрузка и соблюдается порядок.

С ключом. Все записи с одним ключом гарантировано идут в одну партицию (murmur2(key) % partitions). Это полезно, когда важна последовательность по атрибуту: события одного пользователя или одного заказа должны идти по порядку.

Без ключа. До версии 2.3 был round-robin, записи раскидывались поровну, но это давало мелкие неэффективные батчи. С Kafka 2.4 появился sticky механизм: продюсер набивает батч для одной случайной партиции целиком, потом переключается. На дистанции распределение ровное, но на коротком периоде возможны всплески в одну партицию.

Что изменилось в Kafka 3.3+ и 4.0. Старые классы DefaultPartitioner и UniformStickyPartitioner объявили устаревшими (KIP-794) и удалили в 4.0. Теперь дефолтная логика встроена прямо в KafkaProducer (partitioner.class равен null) и включает строго равномерный sticky-партиционер, умеющий подстраиваться под профиль нагрузки: продюсер видит, какие брокеры отвечают медленнее, и льет в них меньше сообщений, выравнивая задержки.

Перегретые партиции. Идеальная картина —  ровный трафик во все партиции без локальных всплесков, но на практике бывают перекосы (skew). Один ключ встречается намного чаще, его партиция становится узким местом, и весь топик упирается в самую медленную партицию. Даже без ключей при рваном потоке часть батчей уходит полными, а часть недозаполненными по таймеру linger.ms.

Батчинг и компрессия

Производительность Kafka во многом держится на отправке данных крупными последовательными наборами сообщений (batch).

Батч против задержки. Собирая сообщения в один запрос, мы можем сильно сократить накладные расходы. Однако этот батч должен еще заполниться, что занимает какое-то время. Этим процессом управляет linger.ms, время удержания открытого батча в надежде доложить туда еще. При linger.ms=0 шлем как можно скорее, минимум задержки на сообщение, но батчи меньше. Больше linger.ms — продюсер ждет и собирает пачку крупнее, эффективность растет ценой небольшой добавки к задержке.

Сменился дефолт в Kafka 4.0. Раньше linger.ms по умолчанию был 0, теперь 5 мс. Выигрыш от крупных батчей обычно дает сопоставимую или даже меньшую итоговую задержку, так что небольшое ожидание окупается. Это один из первых параметров, на который стоит обратить внимание при первичной настройке клиентского приложения.

Размер батча. batch.size — это верхняя граница одного батча в байтах. Крупный размер повышает эффективность использовать сети, однако при слабой нагрузке батч не наберется и уйдет по linger.ms, зря заняв память.

Алгоритмы сжатия. Тут выбор очень зависит от задач, которые стоят перед приложением и Kafka. gzip сжимает лучше всех, но самый прожорливый по CPU и задержке. lz4 и snappy быстрые, но сжатие хуже. zstd пытается совместить высокое сжатие с приемлемой скоростью, но требует дополнительной настройки. Универсального ответа нет, выбор зависит от приоритетов, доступных ресурсов и допустимым задержкам.

ProduceRequest: in-flight и узкие места

Мультиплексирование и in-flight. По одному соединению продюсер шлет несколько запросов, не дожидаясь ответа на предыдущий. В каждом ProduceRequest лежит correlation_id, брокер копирует его в ответ, и продюсер понимает, на какой запрос пришел ответ. Число висящих запросов на соединение ограничено max.in.flight.requests.per.connection (по умолчанию 5). Значение 1 переводит продюсер в синхронный режим: отправил, дождался ACK, отправил следующий.

Зачем это знать. Асинхронная отправка создает риск нарушения порядка записи при сбоях: один батч не получил ACK и отправляется повторно, а следующие за ним уже дошли, в итоге на брокере другой порядок сообщений, а также создается риск появления дублей. Сейчас проблему решает идемпотентный продюсер: при включенной идемпотентности брокер может отловить задублированное сообщение, используя дополнительные механизмы.

Узкое место, сеть или диск. Основная задержка продюсера — это время репликации. При acks=all лидер не только пишет к себе, но и ждет подтверждения от всех реплик ISR. Фактически клиент ждет самый медленный из реплицирующих брокеров, сетевая задержка множится на глубину ISR. При acks=1 лидер отвечает сразу после локальной записи, и задержка сводится к одному round-trip плюс обработка.

Лимит размера сообщения

По умолчанию Kafka ограничивает размер сообщения примерно мегабайтом (message.max.bytes на брокере). Превысили, получили RecordTooLargeException. Лимит выбран не случайно, мелкие сообщения эффективнее буферизуются, упаковываются в пачки и реплицируются, а крупные бьют по всему тракту:

  • Падает пропускная способность, растет задержка.

  • Растет нагрузка на память и GC. Буферы под крупные сообщения выделяются на всех узлах, от продюсера до потребителя.

  • Тормозит репликация и восстановление. Реплицировать и воспроизводить большие записи дороже, перезапуск брокера затягивается.

  • Дороже хранение. С учетом фактора репликации одно тяжелое сообщение занимает несколько своих размеров.

Эти лимиты каскадные и должны быть согласованы на всем тракте. Хотите отправлять большие сообщения — поднимайте синхронно message.max.bytes на брокерах, replica.fetch.max.bytes на репликах, max.request.size у продюсеров и fetch.max.bytes с max.partition.fetch.bytes у потребителей. Некорректная настройка опасна: продюсер отправит большое сообщение, а потребитель не прочитает из-за маленького fetch.max.bytes, и чтение остановится.

4. Надежность записи: acks, ISR и min.insync.replicas

Как мы сказали ранее, продюсер ждет от брокера подтверждения при отправке запроса на запись. Разберем, чем регулируется надежность и где проходит граница, за которой Kafka скорее откажет в записи, чем потеряет данные.

Три уровня acks

Параметр acks задает, сколько реплик должны принять сообщение, прежде чем запись считается успешной.

  • acks=0 (fire-and-forget). Продюсер не ждет ничего, сообщение считается отправленным сразу после передачи в сокет. Минимум задержки, максимум пропускной способности, однако ноль гарантий. Потерялся пакет или упал лидер в момент приема, данные исчезли, продюсер не узнает;

  • acks=1 (подтверждение лидера). Продюсер ждет ACK только от лидера. Лидер пишет к себе и сразу отвечает. Риск —  упал лидер сразу после ответа, не успев реплицировать, новый лидер из ISR этой записи не получит, продюсер думает, что все хорошо, а данных нет. Компромисс скорости при редких потерях;

  • acks=all (подтверждение всеми ISR). Лидер ждет, пока запись не подтвердят все синхронные реплики. Максимальная гарантия: упал лидер, новый выберется из реплик, у которых эта запись уже есть, риск потерять сообщения минимален.

Параметр min.insync.replicas

Сам по себе acks=all еще не гарантирует запись на несколько брокеров, он гарантирует запись на все текущие ISR, а их может остаться и одна (только лидер). Минимальный кворум задает параметр min.insync.replicas (может быть задан на уровне кластера или на уровне конкретного топика), сколько реплик ISR обязаны подтвердить запись.

Частая ошибка выставить min.insync.replicas равным фактору репликации, например, 3 при RF=3. Кажется, что так надежнее всего. На деле для успешной записи теперь нужны все три реплики сразу. Падает один любой брокер, и запись в партицию полностью прекращается с ошибкой NotEnoughReplicas, хотя два из трех узлов живы. Отказоустойчивый кластер превратился в систему, которая не выдерживает потери даже одного узла.

Как падения и сеть превращаются в ошибки записи

В норме Kafka сохраняет данные, пока жива хотя бы одна синхронная реплика. При единичном сбое система умеет переключаться самостоятельно —  текущий лидер партиции упал по тем или иным причинам, новый лидер выбирается из ISR, тем самым уже подтвержденные сообщения не теряются. Для продюсера такой сбой на стороне кластера проявляется в виде временной недоступности или роста задержки.

У гарантий надежности есть границы, которые важно понимать. Набор RF=3, min.insync.replicas=2, acks=all переживает отказ одного брокера без потерь, однако если сбой оказался больше заложенного запаса, например, упали две реплики из трех, то Kafka перестанет принимать новые сообщения. Баланс производительности и надежности во многом достигается с помощью параметров acks и min.insync.replicas: снижаете порог, повышаете шанс писать при крупных сбоях, но рискуете неподтвержденными данными.

5. Дубли и откуда они берутся

Мы научились не терять сообщение, однако у надежности есть и обратная сторона:  в распределенных системах гарантии надежности, как правило, оплачиваются избыточностью. При сбое Kafka скорее продублирует сообщение, чем потеряет его. Разберем источник дублей и как Kafka пытается с ними справляться.

Главный источник дублей, потерянный ответ

Основные причины дублей связаны с неопределенностью. Например, продюсер отправил сообщение, но из-за сетевого сбоя или падения самого брокера никакого ответа не получил. Если сбой произошел до записи в лог, то дубля не будет, но если же сообщение все-таки попало в лог, а ACK потерялся, то продюсер попадает в состояние неопределенности. У него есть два пути:  переотправить и рискнуть дублем, либо ничего не делать и рискнуть потерей сообщений.

Сценарий повторной отправки: Продюсер не получает ACK и поэтому повторно шлет то же самое сообщение, добиваясь успешного подтверждения доставки. Однако первоначальная попытка уже была зафиксирована на стороне кластера, поэтому повтор приводит к двойной записи в журнал Kafka – появляется дубль
Сценарий повторной отправки: Продюсер не получает ACK и поэтому повторно шлет то же самое сообщение, добиваясь успешного подтверждения доставки. Однако первоначальная попытка уже была зафиксирована на стороне кластера, поэтому повтор приводит к двойной записи в журнал Kafka – появляется дубль

Даже в исправном кластере мелкие сбои случаются постоянно, поэтому retries встроены в клиенты по умолчанию.

Как таймауты формируют поведение

Частоту повторной отправки запроса задают три основных параметра, и думать о них надо как о единой системе.

  • request.timeout.ms. Сколько ждать ответа на одну попытку. Истек, попытка неудачна, дальше либо ретрай, либо отказ. Если значение слишком мало, то продюсер паникует рано: брокер еще обрабатывает запрос, а продюсер уже шлет повторный. Слишком большой —  продюсер терпеливый, дольше реагирует на реальные проблемы;

  • delivery.timeout.ms. Общее время на доставку сообщений, со всеми попытками и паузами. При его истечении запись считается окончательно неуспешной. Разумно держать delivery.timeout.ms не меньше request.timeout.ms × (retries + 1);

  • retry.backoff.ms. Пауза между попытками (по умолчанию — 100 мс). Не дает завалить брокер множественными повторами.

Тут как и везде, важен баланс, небольшие значения позволяют быстрее выявлять проблемы, но порождают шквал лишних запросов, трафика и дублей. Большие значения - снижают риск дублей ценой большой задержки при реальных проблемах.

Идемпотентный продюсер: пытаемся защититься от дублей

В Kafka с версии 0.11 появляется идемпотентность продюсера. Механика простая —  при подключении продюсер получает уникальный ProducerId (PID), а его сообщения нумеруются по порядку внутри партиции (sequence number). Брокер помнит максимальный успешно записанный номер для каждого продюсера. Если пришел повтор с уже записанным номером, то брокер ловит дубль и второй раз не пишет. Если номер перепрыгнул последовательность, то это нарушение порядка, и брокер отклоняет запись с OutOfOrderSequence.

Итог: повторы из-за сетевых сбоев больше не дают дублей, порядок сохраняется, потерь нет, и все это при max.in.flight до 5.

Важное про дефолты Kafka 3.0+. Раньше идемпотентность требовало явного включения, теперь enable.idempotence=true и acks=all — это значения по умолчанию. То есть из коробки современный продюсер уже защищен от дублей на транспорте и пишет с максимальной гарантией.

Важно понимать, что идемпотентность защищает только передачу сообщений между продюсером и брокером и только в рамках одной сессии продюсера. Если по какой-то причине приложение перезапустилось, то весь механизм начинает работать заново. Также идемпотентность не спасает от дублей на уровне бизнес-логики —  никто не мешает сервису сформировать одно и то же событие дважды. Эти случаи закрывает следующий уровень, транзакции.

6. Транзакции и Exactly-Once

Представим, что с помощью идемпотентности мы устранили проблему дублей на этапе записи в Kafka. Но возьмем типовую задачу потоковой обработки —  сервис читает из входного топика, обрабатывает и пишет результат в выходной топик. Порядок действий критичен —  записали результат, упали до фиксации offset, после рестарта прочитали тот же вход и снова записали результат, на выходе дубль. Без отдельного механизма нельзя гарантировать, что каждое сообщение из входного топика отразится на выходном ровно один раз.

Идемпотентность и транзакции — это основные механизмы гарантий записи продюсера, поэтому закроем тему целиком, прежде чем уходить к потребителю. Операцию OffsetCommit, которую используют транзакции, подробно разберем во второй части, а пока важна сама идея атомарности.

Атомарность результата и offset

Если вы знакомы с ACID, то решение покажется простым:  сделать связку «прочитал, обработал, записал результат, зафиксировал offset» атомарной. Одна транзакция охватывает весь цикл: приложение начинает транзакцию, публикует выходные сообщения, отправляет offset прочитанного входа и либо коммитит все разом, либо откатывает. Закоммитили, и результаты, и новые offset стали видимы. Откатили, не применилось ничего. Это и есть exactly-once, каждый вход обработан ровно один раз.

За транзакции отвечает координатор транзакций. Он назначает ему внутренние ProducerId (PID) и producer epoch, ведет состояние активных транзакций во внутреннем топике __transaction_state и решает, когда транзакция зафиксирована, а когда отменена. Начало и конец транзакций приложение определяет с помощью методов Producer API: beginTransaction, sendOffsetsToTransaction, commitTransaction, abortTransaction.

Идемпотентность тут никуда не девается, транзакционный продюсер автоматически идемпотентен, она встроена в протокол.

Изоляция чтения

В тему транзакций все-таки стоит упомянуть и чтение, а именно уровень изоляции, параметр isolation.level.

  • read_committed. Потребитель видит только зафиксированное. Брокер держит для него границу LSO (Last Stable Offset), до которой в партиции нет висящих транзакций. Сообщения становятся видны после коммита, отмененные транзакции помечаются управляющими записями и не выдаются вовсе;

  • read_uncommitted. Читает все подряд, не дожидаясь коммита, включая незавершенные и даже отмененные транзакции.

Для сквозного exactly-once потребитель на выходе обязан работать в read_committed.

Где exactly-once заканчивается

Главное, что нужно понять. Транзакции Kafka атомарны только внутри одного кластера Kafka: от чтения входных топиков до записи в выходные и фиксации offset в __consumer_offsets. Какие-либо операции вне этого кластера —  гарантия пропадает.

  • Внешние системы. Пишете попутно в БД, дергаете внешний API, сохраняете файл, эти действия в транзакцию Kafka не входят, сбой в них может создать несогласованное состояние.

  • Многоэтапные распределенные конвейеры. Пайплайн через несколько систем или несколько кластеров Kafka —  в таких крупных системах exactly-once не гарантируется, транзакция локальна для кластера.

  • Ошибки восстановления состояния. In-memory счетчики, кеши и прочее несохраненное при рестарте могут дать другой результат повторной обработки.

Вывод прост: Kafka дает надежный exactly-once внутри своего процессинга, а за его пределами идемпотентность и атомарность это забота архитектуры приложения. Внутри Kafka ни дублей, ни потерь при корректном конфиге не возникнет, однако снаружи нужны свои механизмы: идемпотентные ключи, outbox-паттерн, дедупликация на приемнике и прочие.

Продолжение следует

На этом первая часть подходит к концу. Мы прошли путь записи целиком: от того, как клиент вообще находит нужный брокер, до батчинга, гарантий доставки, идемпотентности и транзакций. Если коротко —  теперь понятно, что происходит между producer.send(...) и попаданием сообщения в лог.

Но это только половина истории. Во второй части развернемся на сто восемьдесят градусов и посмотрим на чтение: как работает pull-модель и long polling, откуда на самом деле берется consumer lag, как живет группа потребителей и почему ребаланс —  это дорого, как фиксируются offset и как квоты умеют незаметно замедлить все, не породив ни одной ошибки в логах.