Kafka: как спроектировать для production — Часть 8
Понимать устройство Kafka недостаточно. Это как знать, как работает двигатель, но всё равно пытаться возить шкаф на самокате.
Большинство проблем возникает не потому, что Kafka работает неправильно, а потому, что на этапе проектирования были неверно выбраны:
- структура
Topics; - количество
Partition; - ключ сообщения;
- срок хранения данных;
- стратегия обработки ошибок;
- модель совместимости сообщений.
Ошибки могут долго оставаться незаметными. Система работает на тестовой нагрузке, успешно проходит запуск, а через несколько месяцев одна Partition оказывается перегруженной, Consumer начинают отставать, а изменение формата события ломает сразу несколько сервисов. Формально всё ещё «зелёное», но где-то в мониторинге уже начинается тихая драма.
Рассмотрим основные решения, которые стоит принять до запуска Kafka в production. Лучше до запуска, чем в пятницу вечером под звук растущего Consumer Lag.
Эта часть является практическим дополнением. В исходной статье подробно рассматриваются архитектура Kafka, Partitions, Consumer Groups, append-only log, репликация и высокая производительность, но рекомендации по проектированию Topics и production-эксплуатации отдельно не разбираются.
Сначала определите назначение события
До создания Topic необходимо ответить на главный вопрос:
Что именно передаётся через Kafka — команда, факт или техническое сообщение?
Рассмотрим три варианта.
Команда
Команда просит конкретный сервис выполнить действие.
GenerateInvoice
SendEmail
ReserveInventory
У команды обычно есть ожидаемый исполнитель.
Например:
Order Service
|
v
GenerateInvoice
|
v
Billing Service
Такие сообщения ближе к модели очереди задач.
Событие
Событие фиксирует факт, который уже произошёл.
OrderCreated
PaymentCompleted
UserRegistered
У события нет единственного получателя.
На него могут подписаться:
- аналитика;
- уведомления;
- антифрод;
- биллинг;
- CRM;
- хранилище данных.
OrderCreated
|
+--------------+--------------+
| | |
v v v
Analytics Notification Fraud
Именно события лучше всего соответствуют Event-Driven архитектуре.
Техническое сообщение
Иногда Kafka используется для внутренних инфраструктурных потоков.
Например:
application-logs
database-changes
metrics-events
audit-records
У таких Topics могут быть другие требования:
- большой объём;
- короткий
Retention; - высокая степень сжатия;
- множество
Partition; - отсутствие строгого порядка между объектами.
Не создавайте один Topic для всего
Одна из первых ошибок — отправлять все события системы в общий Topic. Это удобно примерно до того момента, пока общий Topic не превращается в свалку с красивым названием events.
events
UserRegistered
OrderCreated
PaymentFailed
EmailSent
ProductUpdated
На первый взгляд это удобно.
Producer используют одну точку записи, Consumer фильтруют сообщения по полю type.
Но со временем появляются проблемы:
- разные события требуют разного
Retention; - нагрузка распределяется неравномерно;
- права доступа становятся слишком широкими;
- изменение одного формата влияет на множество
Consumer; - невозможно независимо масштабировать отдельные потоки;
- сложнее отслеживать отставание
Consumer.
Лучше разделять Topics по логическому назначению.
users.events
orders.events
payments.events
notifications.commands
Но Topic на каждый тип события — тоже не всегда хорошо
Противоположная крайность:
order-created
order-paid
order-cancelled
order-shipped
order-delivered
При большом количестве сущностей Kafka-кластер быстро обрастает сотнями или тысячами Topics.
Это усложняет:
- администрирование;
- мониторинг;
- настройку прав;
- управление конфигурацией;
- поддержку схем.
Нередко разумнее объединить события одной бизнес-сущности.
orders.events
А тип события хранить внутри сообщения:
{
"eventType": "OrderPaid",
"eventId": "5c72a8f1",
"orderId": "48291",
"occurredAt": "2026-08-06T09:30:00Z",
"data": {
"amount": 15990,
"currency": "RUB"
}
}
Как именовать Topics
Имя должно помогать понять назначение потока без чтения документации.
Можно использовать структуру:
<domain>.<entity>.<message-type>
Например:
commerce.orders.events
commerce.payments.events
identity.users.events
notifications.email.commands
Если в компании несколько окружений, окружение иногда добавляют в имя:
prod.commerce.orders.events
stage.commerce.orders.events
Однако чаще production и staging размещают в отдельных кластерах или изолируют другими средствами. Тогда окружение в названии только создаёт лишний шум.
Главное правило:
Соглашение должно быть единым для всей компании.
Плохой вариант:
orders
payment_events_v2
UserUpdateTopic
email-queue-prod-new
Хороший вариант:
commerce.orders.events
commerce.payments.events
identity.users.events
notifications.email.commands
Как выбрать Key
Key определяет, в какую Partition попадёт сообщение.
Следовательно, он влияет сразу на три свойства:
- порядок;
- распределение нагрузки;
- возможность масштабирования.
Выбор ключа — одно из важнейших архитектурных решений.
Пример с заказами
Для событий заказа естественным ключом будет orderId.
{
"key": "order-48291",
"eventType": "OrderPaid"
}
Все события одного заказа попадут в одну Partition:
OrderCreated
PaymentStarted
PaymentCompleted
DeliveryStarted
OrderDelivered
Порядок будет сохранён.
Пример с пользователями
Для потока действий пользователя ключом может быть userId.
user-1732
Тогда Kafka сохранит порядок событий одного пользователя:
UserRegistered
EmailChanged
PasswordChanged
UserBlocked
При этом события разных пользователей смогут обрабатываться параллельно.
Плохой ключ: константа
Представим, что у всех сообщений одинаковый Key:
key = orders
Все данные попадут в одну Partition.
Partition 0: 100% нагрузки
Partition 1: 0%
Partition 2: 0%
Kafka-кластер может состоять из десятков Broker, но вся нагрузка всё равно будет проходить через один из них.
Так возникает hot partition.
Плохой ключ: поле с малым количеством значений
Допустим, ключом выбран статус заказа:
created
paid
cancelled
Значений всего три.
При десяти Partition большая часть из них может вообще не использоваться.
Кроме того, значение created будет встречаться гораздо чаще остальных и создаст перекос нагрузки.
Плохой ключ: текущее время
Иногда для равномерного распределения используют Timestamp.
key = 2026-08-06T09:30:01
Нагрузка действительно может распределиться.
Но события одного заказа начнут попадать в разные Partition.
OrderCreated -> Partition 1
PaymentStarted -> Partition 4
PaymentSuccess -> Partition 0
Порядок между ними больше не гарантируется.
Как избежать горячих Partition
Hot partition появляется, когда один ключ встречается значительно чаще остальных.
Например, Topic хранит события продавцов маркетплейса.
Ключ:
sellerId
Один крупный продавец создаёт половину всех событий.
В результате:
seller-1 -> Partition 2 -> 50% трафика
остальные продавцы -> другие Partition
Возможные решения.
Использовать составной ключ
sellerId + shard
Например:
seller-1-0
seller-1-1
seller-1-2
Это распределит события крупного продавца между несколькими Partition.
Но глобальный порядок всех событий продавца будет потерян.
Порядок сохранится только внутри каждого shard.
Выбрать более гранулярный объект
Вместо:
sellerId
можно использовать:
orderId
Если порядок нужен только внутри заказа, это будет более правильный ключ.
Разделить поток
Особенно крупные источники иногда выносят в отдельный Topic.
orders.events
orders.large-sellers.events
Но такое решение увеличивает сложность и требует веских причин.
Сколько Partition создавать
Универсального числа не существует.
Количество зависит от:
- объёма данных;
- требуемой пропускной способности;
- количества
Consumer; - требований к порядку;
- размера сообщений;
- производительности
Broker; - перспектив роста нагрузки.
Важный принцип:
Максимальный параллелизм одной Consumer Group ограничен количеством Partition.
Если в Topic четыре Partition:
P0
P1
P2
P3
эффективно работать одновременно смогут максимум четыре Consumer.
P0 -> Consumer A
P1 -> Consumer B
P2 -> Consumer C
P3 -> Consumer D
Пятый Consumer будет простаивать.
Почему нельзя просто создать тысячу Partition
Каждая Partition требует ресурсов:
- файловых дескрипторов;
- памяти;
- метаданных;
- сетевых соединений;
- операций репликации;
- времени на восстановление;
- времени на
Consumer Rebalance.
Чем больше Partition, тем сложнее кластеру:
- выбирать лидеров;
- перераспределять данные;
- восстанавливаться после падения
Broker; - выполнять
Rebalance Consumer Group.
Поэтому принцип «чем больше, тем лучше» здесь не работает.
Простой подход к расчёту
Допустим, один Consumer обрабатывает:
5 000 сообщений в секунду
Ожидаемая нагрузка:
40 000 сообщений в секунду
Минимально потребуется:
40 000 / 5 000 = 8 Consumer
Следовательно, нужно как минимум восемь Partition.
Но обычно оставляют запас на рост:
12 или 16 Partition
Это не строгая формула, а отправная точка для нагрузочного тестирования.
Планируйте рост заранее
Увеличить количество Partition можно.
Уменьшить — значительно сложнее.
Однако после увеличения меняется результат маршрутизации по ключу:
hash(key) % partition_count
Было:
hash(key) % 4
Стало:
hash(key) % 8
Некоторые новые сообщения начнут попадать в другие Partition.
Из-за этого порядок старых и новых событий одного ключа может нарушиться.
Поэтому количество Partition следует выбирать с учётом будущего роста.
Retention: сколько хранить сообщения
Kafka не удаляет сообщение после того, как его прочитал Consumer.
Удаление определяется политикой хранения.
Основные варианты:
- хранение по времени;
- хранение по объёму;
- Log Compaction;
- сочетание нескольких политик.
Хранение по времени
Например:
Retention = 7 дней
Kafka удаляет сегменты журнала старше этого срока.
Такой режим подходит для:
- технических событий;
- логов;
- промежуточных потоков;
- данных, которые можно восстановить из основной системы.
Хранение по размеру
Можно ограничить максимальный объём Partition.
Retention Size = 100 GB
Когда лимит превышен, старые сегменты удаляются.
Важно понимать, что ограничение обычно применяется к каждой Partition отдельно.
Если Topic состоит из десяти Partition, общий объём может быть значительно больше одного установленного лимита.
Log Compaction
Иногда необходимо хранить не всю историю, а последнее состояние каждого ключа.
Представим поток:
user-1 -> name=Ivan
user-2 -> name=Anna
user-1 -> name=Sergey
После compaction Kafka может сохранить:
user-1 -> name=Sergey
user-2 -> name=Anna
Предыдущая версия user-1 со временем удаляется.
Такой подход полезен для:
- CDC;
- потоков текущего состояния;
- восстановления локального кэша;
- таблиц конфигурации;
- профилей пользователей.
Retention — часть бизнес-требований
Нельзя выбирать Retention случайно.
Нужно ответить:
- сколько времени
Consumerможет быть недоступен; - нужно ли переигрывать историю;
- есть ли данные в другом источнике;
- существуют ли юридические ограничения;
- сколько дискового пространства доступно;
- насколько быстро растёт
Topic.
Например, если Consumer может не работать три дня, а Retention равен одному дню, после восстановления он уже не сможет получить пропущенные события.
Проектируйте формат события
Kafka передаёт байты, но сервисам нужен контракт.
Плохой вариант:
{
"id": 123,
"value": "something"
}
Непонятно:
- что такое
id; - что означает
value; - какая версия сообщения используется;
- когда событие произошло;
- как обнаруживать дубликаты.
Лучше использовать предсказуемую структуру.
{
"eventId": "0ae2342d-105d-45ad-8d91-441b7115afbc",
"eventType": "OrderCreated",
"eventVersion": 1,
"occurredAt": "2026-08-06T09:30:00Z",
"producer": "order-service",
"data": {
"orderId": "48291",
"userId": "1732",
"total": 15990,
"currency": "RUB"
}
}
Зачем нужен Event ID
При модели at-least-once Consumer может получить одно сообщение повторно.
Уникальный идентификатор позволяет проверить, было ли событие уже обработано.
eventId = 0ae2342d
Consumer сохраняет обработанные идентификаторы или использует их как уникальный ключ бизнес-операции.
Зачем нужна версия события
Формат сообщений со временем меняется.
Первая версия:
{
"orderId": "48291",
"total": 15990
}
Вторая версия:
{
"orderId": "48291",
"amount": {
"value": 15990,
"currency": "RUB"
}
}
Старые Consumer могут не понимать новый формат.
Поле версии помогает явно управлять совместимостью:
eventVersion = 2
Не меняйте смысл существующего поля
Допустим, поле price сначала хранит рубли:
{
"price": 1500
}
После обновления оно начинает хранить копейки:
{
"price": 150000
}
Формат остался числом, но смысл изменился.
Это особенно опасно: Consumer не упадут с ошибкой, а начнут незаметно обрабатывать неправильные данные.
Лучше добавить новое поле:
{
"amountMinor": 150000,
"currency": "RUB"
}
Предпочитайте совместимые изменения
Обычно безопаснее:
- добавлять необязательные поля;
- задавать значения по умолчанию;
- не удалять используемые поля;
- не менять тип поля;
- не переиспользовать имя поля с другим смыслом.
Опасные изменения:
string -> integer
optional -> required
поле удалено
значение изменило смысл
Обрабатывайте ошибки отдельно
Consumer не всегда может обработать сообщение.
Причины:
- некорректный формат;
- отсутствующее обязательное поле;
- недоступная база данных;
- ошибка внешнего API;
- нарушение бизнес-правила.
Просто бесконечно повторять одно сообщение опасно.
Offset 100
ошибка
retry
ошибка
retry
ошибка
Одна проблемная запись может остановить всю Partition.
Retry Topic
Временные ошибки можно обрабатывать через отдельный Topic.
orders.events
|
v
Consumer
|
| temporary error
v
orders.retry.1m
Позже сообщение возвращается на повторную обработку.
Можно использовать несколько уровней:
orders.retry.1m
orders.retry.10m
orders.retry.1h
Dead Letter Topic
Если сообщение не удалось обработать после нескольких попыток, его помещают в Dead Letter Topic.
orders.dead-letter
Вместе с исходным сообщением полезно сохранить:
- текст ошибки;
- время последней попытки;
- количество повторов;
- имя
Consumer; - исходный
Topic; Partition;Offset.
{
"sourceTopic": "commerce.orders.events",
"sourcePartition": 3,
"sourceOffset": 15892,
"attempts": 5,
"error": "customer not found",
"failedAt": "2026-08-06T10:15:00Z",
"originalMessage": {
"eventType": "OrderCreated"
}
}
Dead Letter Topic не решает ошибку автоматически.
Он лишь не позволяет одному проблемному сообщению заблокировать весь поток.
Consumer должен быть идемпотентным
При at-least-once одно событие может прийти повторно.
Допустим, Consumer обрабатывает платёж.
Плохая реализация:
Получили PaymentCompleted
|
v
Начислили бонусы
Если сообщение придёт дважды, бонусы начислятся повторно.
Правильнее использовать eventId или идентификатор платежа.
INSERT INTO processed_events (event_id)
VALUES ('0ae2342d')
ON CONFLICT DO NOTHING;
А бизнес-операцию выполнять только для нового идентификатора.
Следите за Consumer Lag
Consumer Lag показывает, насколько Consumer отстаёт от конца Partition.
Последний Offset в Partition: 100000
Offset Consumer Group: 95000
Lag: 5000
Сам по себе lag не всегда означает проблему.
Если Consumer быстро догоняет поток, временное отставание допустимо.
Опасна ситуация, когда lag постоянно растёт:
10:00 -> 5 000
10:05 -> 20 000
10:10 -> 50 000
10:15 -> 100 000
Это значит, что Consumer обрабатывает данные медленнее, чем Producer их создаёт.
Что проверять при росте Lag
Возможные причины:
- недостаточно
Consumer; - недостаточно
Partition; - медленная база данных;
- внешний API отвечает долго;
- слишком тяжёлая бизнес-логика;
- большие сообщения;
- частые паузы сборщика мусора;
- одна горячая
Partition; Consumerпостоянно выполняют rebalance.
Нельзя просто смотреть на общий lag группы.
Нужно проверять отставание по каждой Partition.
Partition 0: Lag 100
Partition 1: Lag 150
Partition 2: Lag 95000
В таком примере проблема, скорее всего, связана с Partition 2.
Избегайте тяжёлой обработки внутри Poll-цикла
Consumer должен регулярно обращаться к Kafka.
Если обработка одного Batch занимает слишком долго, Kafka может решить, что Consumer перестал отвечать.
Это приведёт к rebalance:
Consumer A долго обрабатывает сообщения
|
v
Kafka считает его недоступным
|
v
Partition передаётся Consumer B
|
v
Consumer A возвращается
|
v
Новый Rebalance
Частые rebalance снижают производительность и могут увеличивать количество повторно обработанных сообщений.
Не отправляйте огромные сообщения без необходимости
Kafka хорошо работает с большим количеством относительно небольших сообщений.
Если через неё начинают передавать файлы размером в десятки или сотни мегабайт, появляются проблемы:
- возрастает задержка;
- растёт потребление памяти;
Batchстановится менее эффективным;- увеличивается сетевой трафик;
- репликация занимает больше времени;
- один
Recordможет заблокировать обработкуPartition.
Для больших файлов часто лучше использовать объектное хранилище.
S3 / Object Storage
|
v
Kafka содержит ссылку на файл
Например:
{
"fileId": "report-48291",
"location": "s3://reports/report-48291.pdf",
"checksum": "..."
}
Не используйте Kafka как основную базу данных без причины
Kafka умеет долго хранить события.
Но это не делает её универсальной заменой PostgreSQL, MySQL или ClickHouse.
Kafka неудобна для запросов вида:
SELECT *
FROM orders
WHERE user_id = 1732
AND status = 'paid';
Kafka оптимизирована под последовательную запись и чтение потока, а не под произвольные запросы.
Частая архитектура:
Kafka
|
+--> PostgreSQL
|
+--> ClickHouse
|
+--> Search Index
Kafka доставляет изменения, а специализированные хранилища предоставляют нужные модели чтения.
Документируйте владельца Topic
Для каждого Topic полезно знать:
- какая команда им владеет;
- какие сервисы публикуют данные;
- какие события находятся внутри;
- какой
Keyиспользуется; - сколько длится
Retention; - какая схема сообщения;
- какие гарантии порядка существуют;
- какие
Consumerсчитаются критичными.
Пример карточки Topic:
Topic:
commerce.orders.events
Owner:
Order Platform Team
Key:
orderId
Partitions:
12
Replication Factor:
3
Retention:
14 days
Schema:
OrderEvent v3
Ordering:
Guaranteed per orderId
Без такой документации Kafka быстро превращается в непрозрачную систему, где никто не знает, можно ли удалить старый Topic или изменить формат события.
Контрольный список перед запуском
Перед созданием Topic полезно ответить на следующие вопросы.
Назначение
Это команда или событие?
Кто Producer?
Кто Consumer?
Можно ли добавить новых Consumer позже?
Порядок
Для каких объектов важен порядок?
Какой Key его обеспечивает?
Нет ли риска горячего ключа?
Масштабирование
Какова ожидаемая нагрузка?
Сколько Consumer потребуется?
Хватит ли количества Partition через год?
Хранение
Сколько времени хранить сообщения?
Нужно ли повторное проигрывание?
Что произойдёт, если Consumer не работает несколько дней?
Надёжность
Допустимы ли дубликаты?
Consumer идемпотентен?
Куда попадут необрабатываемые сообщения?
Контракт
Есть ли Event ID?
Есть ли версия схемы?
Совместимы ли изменения со старыми Consumer?
Эксплуатация
Кто владеет Topic?
Какие метрики и алерты настроены?
Как отслеживается Consumer Lag?
Типичная production-схема
Представим сервис заказов.
Order Service
|
| OrderCreated
v
commerce.orders.events
|
+---------------------+
| |
v v
Inventory Group Analytics Group
| |
v v
Reserve Stock Update Metrics
|
| error
v
commerce.orders.retry.1m
|
| repeated failure
v
commerce.orders.dead-letter
Основные решения:
Key:
orderId
Partitions:
12
Replication:
3
Retention:
14 days
Delivery:
at-least-once
Consumer:
idempotent
Error handling:
Retry + Dead Letter Topic
Такая архитектура не гарантирует отсутствие всех проблем, но делает поведение системы предсказуемым.
Главные ошибки при внедрении Kafka
Использовать Kafka только потому, что она популярна
Если системе достаточно HTTP, PostgreSQL или небольшой очереди, Kafka добавит больше эксплуатационных расходов, чем пользы.
Не определять Key
Без осмысленного ключа легко потерять порядок связанных событий.
Создавать слишком мало Partition
Consumer Group не сможет масштабироваться при росте нагрузки.
Создавать слишком много Partition
Кластер будет тратить ресурсы на метаданные, репликацию и rebalance. Партиции не бесплатные, даже если создаются одной командой.
Игнорировать дубликаты
At-least-once без идемпотентности приводит к повторным платежам, уведомлениям и другим бизнес-ошибкам.
Не контролировать схему
Одно несовместимое изменение может одновременно сломать десятки Consumer.
Не следить за Lag
Система формально продолжит работать, но данные будут обрабатываться с задержкой в часы или дни.
Не назначать владельца Topic
Со временем никто не сможет безопасно изменить или удалить поток.
Итог
Kafka начинается не с установки кластера и не с создания Topic.
Она начинается с ответов на архитектурные вопросы:
- что является событием;
- для какого объекта нужен порядок;
- какой
Keyобеспечит этот порядок; - как поток будет масштабироваться;
- сколько времени данные должны храниться;
- как
Consumerобрабатывает дубликаты; - что делать с ошибочными сообщениями;
- как развивать формат данных.
Правильно выбранные Topic, Key и Partition позволяют Kafka масштабироваться годами.
Ошибки в этих решениях могут сделать даже мощный кластер медленным, нестабильным и сложным в сопровождении. Железо многое прощает, но плохую модель данных оно обычно просто делает дороже.
Что дальше?
В следующей части разберём Kafka с точки зрения собеседований и System Design:
- какие вопросы задают чаще всего;
- как объяснить отличие Kafka от RabbitMQ;
- как ответить на вопрос о порядке сообщений;
- как рассуждать о количестве
Partition; - как проектировать Kafka для интернет-магазина;
- какие ошибки в ответах сразу выдают поверхностное понимание.