Кейс: партиционирование по дате переполнило одну партицию Kafka за час

Сервис писал события в топик, ключ партиционирования — дата в формате yyyy-MM-dd. Логика казалась разумной: события одного дня удобно читать вместе, ключ стабильный, партиций хватало с запасом.

Проблема всплыла в день распродажи. Все события за текущий день летели в одну и ту же партицию, потому что ключ у них был одинаковый — сегодняшняя дата. Остальные партиции топика стояли пустыми, а один брокер держал весь входящий поток и не успевал его отдавать consumer'ам. Lag на этой партиции рос, пока остальные consumer'ы группы скучали без работы.

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

Что спросят следом: как выбрать ключ правильно. Ответ — добавить в ключ что-то с высокой кардинальностью и равномерным распределением, например userId или orderId, а дату оставить как поле в самом сообщении, а не в ключе. Ещё спросят, как обнаружить такую проблему до продакшена: смотреть на распределение сообщений по партициям через consumer group lag по каждой партиции отдельно, а не по топику в среднем — усреднённый lag прячет именно такие перекосы.

Ключ партиционирования — это решение о параллелизме, а не о смысловой группировке.

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

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


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