Мне показывали пост-мортем этой истории.
В Кафку шел большой поток данных.
1. Консьюмер вычитывал пачку сообщений и приступал к их обработке.
2. Приложение с консьюмером не успевало обработать всю пачку за время, отведенное консьюмеру, чтобы успеть запросить новую пачку, и закоммитить оффсет.
3. Координатор, не получив за отведенное время запроса новой пачки, выкидывал бедного консьюмера из группы, а партицию переназначал другому консьюмеру.
4. Другой консьюмер вычитывал пачку сообщений и приступал к их обработке... и далее всё повторялось (см. п. 1)
То есть оффсеты не коммитились, консьюмеры по кругу вычитывали одни и те же сообщения, а количество необработанных сообщений в топике все росло и росло 😱
Как это чинили:
⚡️ Уменьшили количество сообщений, которые вычитывает консьюмер за один запрос (
max.poll.records)⚡️ Увеличили время, в течение которого координатор ждет от консьюмера нового запроса прежде чем счесть его умершим и выкинуть из группы (
max.poll.interval.ms)⚡️ Увеличили количество партиций и консьюмеров
🧐 Можно просто назначить консьюмеру топик или набор партиций и чтобы они не переназначались?
Можно. Только если вдруг будут добавляться новые партиции в топик, этот "прибитый гвоздями" консьюмер о них не узнает сам, нужно будет перезапустить приложение. И
group.id все равно надо указать, иначе коммиты оффсетов не будут в Кафке сохраняться.#kafka