Kafka — это не очередь. Это журнал, который никто не сотрёт

Если вы работали с очередями вроде RabbitMQ или SQS, у Kafka есть одна особенность, которая ломает привычную интуицию: прочитанное сообщение никуда не девается. Разберёмся, как это устроено и что из этого следует — на простых примерах, без воды.

Слева — как работает очередь. Справа — как работает Kafka

Почему это не очередь

В очереди сообщение живёт до первого прочтения: воркер забрал — брокер удалил. Удобно для задач «сделай что-то один раз», но плохо подходит, если один и тот же поток данных нужен нескольким системам сразу.

Kafka устроена иначе. Сообщение дописывается в конец файла на диске — такой файл называют логом, или commit log — и хранится там столько, сколько настроено, независимо от того, кто и сколько раз его уже прочитал.

Продюсер только дописывает справа. Потребители читают с любой позиции (offset) и не удаляют прочитанное

Отсюда главный практический эффект: пять разных систем — аналитика, поиск, антифрод, уведомления, аудит — могут читать один и тот же поток независимо, каждая в своём темпе, не мешая друг другу.
▫️O(1) — чтение и запись это последовательный доступ к диску, не поиск по индексу
▫️ — сколько угодно независимых читателей без доп. нагрузки на источник
▫️TTL / ∞ — хранение настраивается: от часов до бессрочного

Из чего состоит кластер

Три термина, без которых дальше не разобраться:

  • Topic — логический поток событий, например dialog-events.
  • Partition — топик физически режется на партиции: независимые упорядоченные логи. Это единица параллелизма.
  • Broker — сервер кластера, хранящий партиции. У каждой партиции один лидер и несколько реплик-фолловеров на других брокерах.
Топик на 3 партиции с replication factor = 2 — у каждой партиции один лидер и реплики на других брокерах

Тут есть два правила, которые часто путают:

Порядок гарантирован только внутри партиции. Между партициями — нет. Нужен строгий порядок между конкретными событиями (например, все события одного диалога) — они должны попадать в одну партицию.

За отказоустойчивость отвечают реплики, а не партиции. Продюсеры и консьюмеры общаются только с лидером. Фолловеры молча копируют его лог. Упал лидер — один из синхронных фолловеров (ISR) становится новым лидером автоматически.

Как не потерять порядок

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

  • Нужен порядок внутри диалога или у одного юзера → задайте key = dialog_id (или user_id).
  • Порядок не важен, важна равномерная нагрузка → не задавайте key — Kafka раскидает сообщения round-robin.

Число партиций топика нельзя безболезненно уменьшить, а увеличение задним числом рвёт гарантию порядка для уже написанных данных. Планируйте количество партиций с запасом, а не наращивайте реактивно.

Кто и как читает

Консьюмеры объединяются в consumer group по общему group.id. Правило простое: одну партицию внутри группы читает только один консьюмер одновременно — это и есть встроенный механизм масштабирования обработки.

Четыре партиции распределены между двумя консьюмерами одной группы; добавление третьего вызовет rebalance

Партиций больше, чем консьюмеров — лишние простаивают. Число партиций стоит планировать исходя из желаемого параллелизма обработки, а не только из throughput записи.

Офсет, который коммитит консьюмер — это отметка «докуда я дочитал», а не команда на удаление. Тот же топик в любой момент может независимо перечитать другая consumer group.

Насколько надёжна доставка

«Exactly-once» в Kafka — это не волшебная кнопка, а комбинация настроек, за которую вы платите производительностью. По умолчанию включён средний по надёжности вариант.

Три уровня гарантий доставки — от самого быстрого до самого надёжного

Две настройки, которые управляют этим спектром:

  • acks — сколько реплик должны подтвердить запись: 0 (никого не ждать), 1 (только лидер), all (все синхронные реплики).
  • enable.idempotence — брокер присваивает каждому сообщению номер и сам отбрасывает дубли от повторной отправки при сетевых ретраях.

Retention: два способа не хранить всё вечно

По умолчанию (cleanup.policy=delete) сообщения стираются по возрасту или объёму — обычный поток событий, где старое со временем не нужно.

Log compaction (cleanup.policy=compact) — Kafka хранит только последнее значение для каждого ключа. Топик превращается в снэпшот текущего состояния, а не в историю: удобно для восстановления состояния сервиса после падения (event sourcing) или для «текущего профиля пользователя».

Масштабирование в двух словах

Разные оси дают разный эффект, и их легко перепутать:

  • Больше партиций → выше throughput записи и параллелизм чтения. Но больше файловых дескрипторов на брокер и дольше выборы лидера при отказе.
  • Больше реплик (replication factor) → надёжнее, но не быстрее. Клиенты всегда общаются только с лидером — реплики не ускоряют ни чтение, ни запись.

min.insync.replicas — недооценённая настройка: сколько реплик должно быть «в синхроне», чтобы запись с acks=all считалась успешной. При replication.factor=3 и min.insync.replicas=2 кластер переживёт отказ одного брокера без остановки записи; при отказе двух — запись остановится вместо тихой потери данных.

Экосистема вокруг Kafka

  • Kafka Connect — готовые коннекторы для загрузки данных в Kafka и выгрузки из неё (например, Debezium для CDC из Postgres).
  • Kafka Streams — библиотека для потоковой обработки прямо поверх топиков, обычное JVM-приложение.
  • ksqlDB — те же потоковые трансформации, но на SQL, без кода.
  • Schema Registry — версионирование схем сообщений и проверка обратной совместимости для команд, которые пишут/читают один топик.

Когда Kafka уместна, а когда нет

Хорошо подходит, если:

  • нужно, чтобы несколько систем читали один поток независимо;
  • важен большой throughput на обычном железе;
  • нужна возможность переиграть историю (реплей) после бага в обработке.

Плохо подходит, если:

  • нужны SQL-запросы и фильтрация по произвольному полю — Kafka не хранилище, только последовательное чтение;
  • сообщений мало и они очень мелкие — overhead на каждое делает лёгкие очереди вроде SQS эффективнее;
  • нет ресурсов на операционную поддержку — свой кластер это ZooKeeper/KRaft, мониторинг ISR и ребалансировка.
Инструмент Что сильнее у него Когда взять вместо Kafka
Kafka throughput, независимые консьюмеры, реплей
RabbitMQ гибкая маршрутизация, приоритеты Task-очередь: сообщение нужно ровно одному воркеру
Amazon SQS zero-ops, простота в AWS Простой pipeline без своей инфраструктуры
Redis Streams минимальная latency Небольшой объём, лёгкий кейс, Redis уже в стеке
Apache Pulsar мультиарендность Нужно масштабировать storage отдельно от compute

Где Kafka встречается в реальных проектах

  1. Шина событий между микросервисами — вместо синхронных вызовов друг друга.
  2. CDC — Debezium стримит изменения из Postgres/MySQL в реальном времени.
  3. Приём потока в аналитику — события пишутся в Kafka, оттуда попадают в ClickHouse и подобные хранилища.
  4. Log aggregation — центральный буфер логов и метрик перед системами мониторинга.

Минимальный рабочий пример

kafka-topics.sh --create --topic dialog-events \
  --partitions 6 --replication-factor 3 \
  --config retention.ms=2592000000   # 30 дней

# producer.properties — надёжная запись
acks=all
enable.idempotence=true
linger.ms=5   # небольшая задержка ради батчинга — выше throughput
// пишем с ключом dialog_id — порядок гарантирован внутри диалога
producer.send(new ProducerRecord<>("dialog-events", dialogId, eventJson));

// коммитим offset после обработки, а не до — at-least-once
while (true) {
    var records = consumer.poll(Duration.ofMillis(500));
    for (var record : records) process(record);
    consumer.commitSync();
}
Отправить
Поделиться
Твитнуть
Запинить