Kafka: как спроектировать для production — Часть 8

Содержание
Коротко: Production-дизайн Kafka начинается до запуска кластера: нужно выбрать Topic, Key, количество Partition, Retention, формат событий, retry/DLQ и правила идемпотентной обработки.

Понимать устройство 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 для интернет-магазина;
  • какие ошибки в ответах сразу выдают поверхностное понимание.