Кейс: consumer вылетал из группы посреди обработки одного батча

Сервис читал сообщения из Kafka и обрабатывал каждое с обращением к внешнему API. При деградации API обработка одного poll-батча растягивалась дольше пяти минут — дефолтного max.poll.interval.ms.

Heartbeat уходил исправно: поток отправки heartbeat отдельный от потока обработки, начиная с KIP-62. Но брокер следит не только за heartbeat — он ждёт, что consumer вызовет poll() вовремя. Не вызвал — группа считает его зависшим и запускает rebalance.

Партиции переехали к другому consumer'у до того, как первый закоммитил offset уже обработанных сообщений. Новый владелец начал читать с последнего закоммиченного offset — часть сообщений обработалась второй раз.

Что спросят следом: как чинить. Вариант первый — снизить max.poll.records, чтобы батч обрабатывался быстрее интервала. Вариант второй — увеличить max.poll.interval.ms, если обработка объективно долгая. Вариант третий — вынести тяжёлую обработку в отдельный пул потоков, а poll-loop оставить быстрым и коммитить offset после подтверждения из пула.

Rebalance из-за таймаута — это не баг Kafka, это защита от зависшего consumer'а, которая срабатывает и на здоровом, но медленном коде.

Тренажёр: 600 вопросов, мок с таймером, план повторов

senior·base — что спрашивают на самом деле


В этом посте были ссылки, но мы их удалили по правилам Сетки