Известно, что Кафка обеспечивает порядок сообщений в рамках партиции. Если мы отправили сообщения в разные партиции одного топика, то гарантии порядка между ними не будет. Но если говорить об одной партиции, может ли нарушиться порядок?
С высоты птичьего полета процесс выглядит так:
Продюсеры отправляют данные -> Брокеры записывают данные на диск -> Консьюмеры забирают данные для обработки.
➡️ Начнем с хранения.
Партиция в Kafka — это последовательный append-only log. Новые сообщения не вставляются в середину и не перезаписывают старые, а просто добавляются в конец файла и так лежат на диске. Брокер не сортирует и не тасует сообщения, поэтому перепутаться непосредственно при хранении не могут.
➡️ Теперь посмотрим на продюсера.
Есть комбинация настроек продюсера, при которой может произойти нарушение порядка:
enable.idempotence=false
retries > 0 //сколько раз продюсер будет переотправлять запрос при ошибке
max.in.flight.requests.per.connection > 1 //сколько запросов продюсер может отправить брокеру, не дожидаясь подтверждения предыдущих
Сценарий такой:
Исходный порядок сообщений: 1, 2, 3, 4, 5, 6.
1. Продюсер отправляет batch A: сообщения 1, 2, 3
2. Не дожидаясь ответа, отправляет batch B: сообщения 4, 5, 6
3. Batch A з-за кратковременного сбоя не записался или потерялось его подтверждение записи
4. Batch B успешно записался в партицию
5. Продюсер ретраит batch A
6. Batch A успешно записывается позже batch B 😱
В логе Kafka сообщения оказываются записаны в порядке: 4, 5, 6, 1, 2, 3
Этот случай прямо описывается в доке Кафки.
Как избежать такой ситуации?
1️⃣ Максимально строгий и простой режим:
max.in.flight.requests.per.connection = 1
Тогда продюсер не будет держать несколько незавершенных запросов одновременно, и сценарий "второй батч записался раньше первого" исчезнет.
➖ Но такая настройка снижает пропускную способность.
2️⃣Включение идемпотентности продюсера и соответствующей ей комбинации настроек:
enable.idempotence = true
retries > 0
max.in.flight.requests.per.connection <= 5
acks = all //сколько реплик партиции должны подтвердить, что получили запись, чтобы продюсер счел отправку успешной
Для версий Кафки с 3.2.0 это настройки по-умолчанию.
Как это работает
Когда идемпотентный продюсер стартует, он получает от брокера уникальный Producer ID. Для каждой партиции продюсер ведёт счётчик sequence number. Эти sequence numbers добавляется в батчи, который продюсер отправляет брокеру. На брокере для пары producerId + partition хранится состояние: батч с каким sequence number был принят последним. Так брокер понимает: “этот батч я уже видел, второй раз писать не надо” или "стоп, этот батч пришел не по порядку, я ещё не видел батч с предыдущим номером".
Так Кафка защищается от дублей из-за ретраев продюсера и сохраняет порядок записи в рамках партиции.
При этом даже идемпотентный продюсер в редких сценариях может отправлять дубли, а консьюмер может вычитывать одно и то же сообщение несколько раз (если упал после обработки, но до коммита оффсета). Поэтому важно следить, чтобы консьюмер был идемпотентным. К сожалению, я наступала на эти грабли, поэтому запомнила не только в теории, но и на практике.
➡️ А может ли порядок сообщений нарушиться уже на стороне консьюмера?
Да, но не из-за самой Кафки, а из-за нашего кода обработки. Например, из-за параллельной обработки или из-за ретраев с отдельным топиком.
1. Консьюмер получил из партиции сообщение 1
2. При обработке сообщения 1 возник временный сбой, из-за чего обработка не удалась, и сообщение 1 было отправлено в отдельный retry-топик
3. Консьюмер получил из партиции сообщения 2, 3
4. Сообщения 2 и 3 были успешно обработаны
5. Из retry-топика было вычитано сообщение 1 и успешно обработано
Фактический порядок обработки стал: 2, 3, 1.
Поэтому если порядок критичен, нужно аккуратно проектировать ретрай-логику.
Таким образом, Кафка гарантирует порядок внутри партиции, но эту гарантию можно сломать настройками продюсера или логикой обработки на стороне консьюмера.
