← Все темы
Kafka
Вопросов: 33
Kafka — распределённый лог событий (distributed event log), он же брокер сообщений для обмена данными между сервисами в реальном времени.
- Модель publish/subscribe: producer'ы пишут события в топики, consumer'ы их читают.
- Данные хранятся на диске и не удаляются после прочтения — их можно перечитывать (replay).
- Горизонтально масштабируется и переживает отказы брокеров.
Зачем на практике: интеграция микросервисов, событийная архитектура, потоковая обработка, сбор логов и метрик, буфер между быстрым producer'ом и медленным consumer'ом.
- Модель publish/subscribe: producer'ы пишут события в топики, consumer'ы их читают.
- Данные хранятся на диске и не удаляются после прочтения — их можно перечитывать (replay).
- Горизонтально масштабируется и переживает отказы брокеров.
Зачем на практике: интеграция микросервисов, событийная архитектура, потоковая обработка, сбор логов и метрик, буфер между быстрым producer'ом и медленным consumer'ом.
Три базовые роли в Kafka:
- Producer — пишет (публикует) сообщения в топик. Это твоё приложение.
- Consumer — подписывается на топик и читает сообщения. Тоже твоё приложение.
- Broker — сервер Kafka: хранит данные и обслуживает запросы producer'ов и consumer'ов. Несколько брокеров образуют кластер.
Коротко: producer и consumer — это твой код, broker — сама Kafka.
- Producer — пишет (публикует) сообщения в топик. Это твоё приложение.
- Consumer — подписывается на топик и читает сообщения. Тоже твоё приложение.
- Broker — сервер Kafka: хранит данные и обслуживает запросы producer'ов и consumer'ов. Несколько брокеров образуют кластер.
Коротко: producer и consumer — это твой код, broker — сама Kafka.
- Топик (topic) — именованный канал сообщений (например
- Партиция (partition) — топик физически разбит на N партиций; это единица параллелизма и хранения.
- Внутри одной партиции сообщения строго упорядочены; между партициями порядка нет.
- Больше партиций → выше параллелизм: в группе одну партицию читает один consumer.
Итог: топик — это логика, партиции — это масштаб и параллелизм.
orders). Producer пишет в топик, consumer читает из топика.- Партиция (partition) — топик физически разбит на N партиций; это единица параллелизма и хранения.
- Внутри одной партиции сообщения строго упорядочены; между партициями порядка нет.
- Больше партиций → выше параллелизм: в группе одну партицию читает один consumer.
Итог: топик — это логика, партиции — это масштаб и параллелизм.
Offset — порядковый номер сообщения внутри партиции (0, 1, 2, …), монотонно растёт.
- Уникален в пределах партиции, а не всего топика.
- Consumer хранит (коммитит) свой offset — «до какого места прочитал», чтобы после перезапуска продолжить, а не начать заново.
- Broker не удаляет сообщение после чтения: несколько групп читают одну партицию независимо, у каждой свой offset.
- Уникален в пределах партиции, а не всего топика.
- Consumer хранит (коммитит) свой offset — «до какого места прочитал», чтобы после перезапуска продолжить, а не начать заново.
- Broker не удаляет сообщение после чтения: несколько групп читают одну партицию независимо, у каждой свой offset.
Сообщение (record) состоит из:
- key — ключ (опционально); определяет партицию и используется в log compaction.
- value — полезная нагрузка (payload), обычно JSON / Avro / Protobuf в виде байтов.
- headers — метаданные (trace-id, тип события и т.п.).
- timestamp — время события или записи.
- Плюс служебное: топик, партиция, offset.
Producer сериализует key/value в байты; Kafka хранит и отдаёт байты и о содержимом ничего не знает.
- key — ключ (опционально); определяет партицию и используется в log compaction.
- value — полезная нагрузка (payload), обычно JSON / Avro / Protobuf в виде байтов.
- headers — метаданные (trace-id, тип события и т.п.).
- timestamp — время события или записи.
- Плюс служебное: топик, партиция, offset.
Producer сериализует key/value в байты; Kafka хранит и отдаёт байты и о содержимом ничего не знает.
- Есть ключ → партиция =
- Нет ключа → сообщения распределяются по партициям равномерно (round-robin / sticky).
- Практика: ключ = сущность, для которой важен порядок (
- Осторожно: «горячий» ключ перегружает одну партицию (перекос), а изменение числа партиций ломает соответствие ключ→партиция.
hash(key) % число_партиций. Одинаковый ключ всегда идёт в одну партицию → сохраняется порядок для этого ключа.- Нет ключа → сообщения распределяются по партициям равномерно (round-robin / sticky).
- Практика: ключ = сущность, для которой важен порядок (
user_id, order_id).- Осторожно: «горячий» ключ перегружает одну партицию (перекос), а изменение числа партиций ломает соответствие ключ→партиция.
- Каждая партиция — append-only лог: новые сообщения дописываются в конец, старые неизменяемы.
- Лог разбит на сегменты (файлы) плюс индексы
- Последовательная запись и чтение с диска + page cache ОС → очень быстро.
- Данные живут независимо от того, прочитаны они или нет: удаляются по retention, а не после чтения.
- Лог разбит на сегменты (файлы) плюс индексы
offset → позиция для быстрого поиска.- Последовательная запись и чтение с диска + page cache ОС → очень быстро.
- Данные живут независимо от того, прочитаны они или нет: удаляются по retention, а не после чтения.
Retention — политика, сколько хранить сообщения перед удалением.
- По времени:
- По размеру:
- Что наступит раньше — то и сработает; удаляются целые старые сегменты.
- Настраивается на уровне топика.
Важно: если consumer lag превысит retention — данные удалятся раньше, чем ты их прочитаешь.
- По времени:
retention.ms (например, 7 дней).- По размеру:
retention.bytes.- Что наступит раньше — то и сработает; удаляются целые старые сегменты.
- Настраивается на уровне топика.
Важно: если consumer lag превысит retention — данные удалятся раньше, чем ты их прочитаешь.
- Compaction (
- Обычный retention удаляет по возрасту/размеру, независимо от ключа.
- Смысл: топик как «снимок текущего состояния» (последний адрес пользователя, актуальный конфиг).
- Удаление ключа — сообщение с этим key и
- Часто используется для changelog- и compacted-топиков.
cleanup.policy=compact) хранит хотя бы последнее значение для каждого ключа, удаляя старые версии того же ключа.- Обычный retention удаляет по возрасту/размеру, независимо от ключа.
- Смысл: топик как «снимок текущего состояния» (последний адрес пользователя, актуальный конфиг).
- Удаление ключа — сообщение с этим key и
value=null (tombstone).- Часто используется для changelog- и compacted-топиков.
- У каждой партиции
- Одна реплика — leader: через неё идут все записи и (обычно) чтения. Остальные — followers, копируют данные с лидера.
- Упал брокер с лидером → новым лидером становится один из синхронных followers, данные не теряются.
- Producer и consumer всегда работают с лидером партиции.
replication.factor копий на разных брокерах.- Одна реплика — leader: через неё идут все записи и (обычно) чтения. Остальные — followers, копируют данные с лидера.
- Упал брокер с лидером → новым лидером становится один из синхронных followers, данные не теряются.
- Producer и consumer всегда работают с лидером партиции.
ISR — набор реплик (лидер + followers), которые «догнали» лидера и не отстают.
- Новым лидером может стать только реплика из ISR → не теряем подтверждённые данные.
- Отставшая реплика выпадает из ISR; догонит — вернётся.
- Работает в связке с
- Новым лидером может стать только реплика из ISR → не теряем подтверждённые данные.
- Отставшая реплика выпадает из ISR; догонит — вернётся.
- Работает в связке с
acks=all и min.insync.replicas.
acks — сколько подтверждений от брокеров ждать, прежде чем считать запись успешной. Это компромисс «надёжность ↔ задержка».
- acks=0 — не ждём вовсе. Максимально быстро, но можно потерять (fire-and-forget).
- acks=1 — ждём только лидера. Лидер упал до репликации → сообщение теряется.
- acks=all (
Для важных данных:
- acks=0 — не ждём вовсе. Максимально быстро, но можно потерять (fire-and-forget).
- acks=1 — ждём только лидера. Лидер упал до репликации → сообщение теряется.
- acks=all (
-1) — ждём лидера и все реплики из ISR. Максимально надёжно.Для важных данных:
acks=all + идемпотентный producer.
min.insync.replicas — минимум синхронных реплик, при котором разрешена запись с
- Реплик в ISR меньше порога → broker отклоняет запись (
Классика durability:
acks=all.- Реплик в ISR меньше порога → broker отклоняет запись (
NotEnoughReplicas), чтобы не принять легко теряемые данные.Классика durability:
replication.factor=3, min.insync.replicas=2, acks=all — переживает отказ одного брокера без потери данных и без остановки записи.
Отказоустойчивость складывается из нескольких механизмов:
- Репликация партиций на разные брокеры.
- Выбор нового лидера из ISR при отказе.
-
- Durable append-only лог на диске переживает перезапуск брокера.
- Consumer после сбоя продолжает с закоммиченного offset.
Итог: ни отказ брокера, ни перезапуск consumer'а не теряют подтверждённые данные.
- Репликация партиций на разные брокеры.
- Выбор нового лидера из ISR при отказе.
-
acks=all + min.insync.replicas → подтверждённое сообщение есть на нескольких репликах.- Durable append-only лог на диске переживает перезапуск брокера.
- Consumer после сбоя продолжает с закоммиченного offset.
Итог: ни отказ брокера, ни перезапуск consumer'а не теряют подтверждённые данные.
Consumer group — группа consumer'ов с общим
- Каждая партиция в группе достаётся ровно одному consumer'у → нагрузка делится.
- Consumer'ов больше, чем партиций → лишние простаивают (параллелизм ограничен числом партиций).
- Разные группы читают один топик независимо, у каждой свои offset (broadcast между группами).
group.id, которые вместе читают топик.- Каждая партиция в группе достаётся ровно одному consumer'у → нагрузка делится.
- Consumer'ов больше, чем партиций → лишние простаивают (параллелизм ограничен числом партиций).
- Разные группы читают один топик независимо, у каждой свои offset (broadcast между группами).
Rebalance — перераспределение партиций между consumer'ами группы.
- Триггеры: consumer вошёл / вышел / упал (перестал слать heartbeat), изменилось число партиций.
- Во время rebalance потребление приостанавливается («stop-the-world») — заметная пауза.
- Частые rebalance — это проблема: долгая обработка превышает
- Смягчают: cooperative rebalancing, static membership, корректные таймауты.
- Триггеры: consumer вошёл / вышел / упал (перестал слать heartbeat), изменилось число партиций.
- Во время rebalance потребление приостанавливается («stop-the-world») — заметная пауза.
- Частые rebalance — это проблема: долгая обработка превышает
max.poll.interval.ms, нестабильные consumer'ы.- Смягчают: cooperative rebalancing, static membership, корректные таймауты.
- Auto-commit (
- Manual commit — коммитишь сам после успешной обработки → at-least-once (дубли возможны, потерь нет).
Практика для надёжности: auto-commit выключить, «обработал → закоммитил», а саму обработку сделать идемпотентной.
enable.auto.commit=true) — Kafka периодически коммитит offset сама. Просто, но при сбое возможны потери или дубли (коммит до обработки → потеря).- Manual commit — коммитишь сам после успешной обработки → at-least-once (дубли возможны, потерь нет).
Практика для надёжности: auto-commit выключить, «обработал → закоммитил», а саму обработку сделать идемпотентной.
Lag = (последний offset в партиции) − (закоммиченный offset consumer'а). Насколько мы отстаём.
- Растёт → consumer не успевает за producer'ом (медленная обработка, мало партиций, упал consumer).
- Ключевая метрика здоровья: большой lag → задержки и риск не успеть прочитать до retention.
- Смотрят через
- Лечат: больше consumer'ов (до числа партиций), ускорить обработку, батчинг.
- Растёт → consumer не успевает за producer'ом (медленная обработка, мало партиций, упал consumer).
- Ключевая метрика здоровья: большой lag → задержки и риск не успеть прочитать до retention.
- Смотрят через
kafka-consumer-groups, Burrow, Prometheus/Grafana.- Лечат: больше consumer'ов (до числа партиций), ускорить обработку, батчинг.
- at-most-once — максимум один раз: возможна потеря, дублей нет (коммит offset до обработки).
- at-least-once — минимум один раз: потерь нет, возможны дубли (коммит после обработки). Дефолт на практике.
- exactly-once — ровно один раз: ни потерь, ни дублей.
Exactly-once достигается идемпотентным producer'ом + транзакциями (EOS) либо at-least-once + идемпотентной обработкой на стороне consumer'а.
- at-least-once — минимум один раз: потерь нет, возможны дубли (коммит после обработки). Дефолт на практике.
- exactly-once — ровно один раз: ни потерь, ни дублей.
Exactly-once достигается идемпотентным producer'ом + транзакциями (EOS) либо at-least-once + идемпотентной обработкой на стороне consumer'а.
- При retry (таймаут, сбой сети) producer может отправить одно сообщение дважды → дубль в партиции.
- Идемпотентный producer (
- Гарантирует «не более одного раза» на запись в партицию и сохраняет порядок при retry.
- Сейчас включён по умолчанию (с
- Идемпотентный producer (
enable.idempotence=true) даёт каждому producer'у PID и порядковые номера; broker отбрасывает повторы.- Гарантирует «не более одного раза» на запись в партицию и сохраняет порядок при retry.
- Сейчас включён по умолчанию (с
acks=all). Дёшево — держать включённым.
- Транзакции — атомарная запись в несколько партиций/топиков плюс коммит offset'ов consumer'а в одной операции. Всё или ничего.
- Основа паттерна consume-process-produce (прочитал → обработал → записал) без дублей = exactly-once.
- Задаётся
- Цена: сложнее и медленнее — включают только там, где реально нужен EOS.
- Основа паттерна consume-process-produce (прочитал → обработал → записал) без дублей = exactly-once.
- Задаётся
transactional.id; consumer читает с isolation.level=read_committed (видит только закоммиченное).- Цена: сложнее и медленнее — включают только там, где реально нужен EOS.
at-least-once означает, что дубли возможны, поэтому обработку делают идемпотентной.
- Уникальный ключ события + проверка «уже обработано» (таблица
-
- Дедуп по
Часто это проще и надёжнее, чем настраивать полноценный exactly-once через транзакции.
- Уникальный ключ события + проверка «уже обработано» (таблица
processed_ids).-
upsert вместо insert; установка значения вместо инкремента (natural idempotency).- Дедуп по
(topic, partition, offset) или по бизнес-ключу.Часто это проще и надёжнее, чем настраивать полноценный exactly-once через транзакции.
- Порядок гарантирован только внутри одной партиции, а не по топику целиком.
- Нужен порядок для сущности → шли её события с одним ключом (все события
- Несколько партиций → глобального порядка нет; это цена параллелизма.
- Идемпотентный producer сохраняет порядок при retry; без него
- Нужен порядок для сущности → шли её события с одним ключом (все события
order_id попадут в одну партицию).- Несколько партиций → глобального порядка нет; это цена параллелизма.
- Идемпотентный producer сохраняет порядок при retry; без него
max.in.flight>1 + retry могут переставить сообщения.
- RabbitMQ — классическая брокер-очередь (smart broker): маршрутизация, сообщение удаляется после ack, push-модель. Хорош для задач/команд, RPC, сложной маршрутизации.
- Kafka — распределённый лог (dumb broker, smart consumer): сообщения хранятся и перечитываются, pull-модель, высокий throughput, replay.
Выбор: событийный поток, аналитика, много подписчиков, replay, большие объёмы → Kafka. Очередь задач, гибкая маршрутизация, низкая латентность на сообщение → RabbitMQ.
- Kafka — распределённый лог (dumb broker, smart consumer): сообщения хранятся и перечитываются, pull-модель, высокий throughput, replay.
Выбор: событийный поток, аналитика, много подписчиков, replay, большие объёмы → Kafka. Очередь задач, гибкая маршрутизация, низкая латентность на сообщение → RabbitMQ.
- Kafka хранит просто байты; о структуре договариваются producer и consumer.
- Форматы: JSON (просто, многословно, без схемы), Avro / Protobuf (компактно, бинарно, со схемой).
- Schema Registry — сервис со схемами; в сообщении едет ID схемы вместо самой схемы.
- Даёт контроль совместимости (backward / forward): producer не сломает consumer'ов несовместимым изменением.
Решает проблему эволюции схемы данных.
- Форматы: JSON (просто, многословно, без схемы), Avro / Protobuf (компактно, бинарно, со схемой).
- Schema Registry — сервис со схемами; в сообщении едет ID схемы вместо самой схемы.
- Даёт контроль совместимости (backward / forward): producer не сломает consumer'ов несовместимым изменением.
Решает проблему эволюции схемы данных.
- Kafka Connect — готовые коннекторы для интеграции без кода: source (БД/файлы → Kafka) и sink (Kafka → БД/ES/S3). Основа CDC (Debezium).
- Kafka Streams — Java-библиотека для потоковой обработки: map/filter/join, агрегации и окна прямо в приложении, без отдельного кластера.
Как бэкендеру: Connect — «перекачать данные», Streams — «обработать поток на лету».
- Kafka Streams — Java-библиотека для потоковой обработки: map/filter/join, агрегации и окна прямо в приложении, без отдельного кластера.
Как бэкендеру: Connect — «перекачать данные», Streams — «обработать поток на лету».
- Раньше Kafka хранила метаданные (брокеры, топики, лидеры) во внешнем ZooKeeper.
- KRaft (KIP-500) — Kafka сама управляет метаданными через встроенный Raft-консенсус, ZooKeeper не нужен.
- Плюсы: проще эксплуатация (одна система), быстрее failover, лучше масштаб по числу партиций.
- В новых версиях KRaft — по умолчанию, ZooKeeper объявлен deprecated.
- KRaft (KIP-500) — Kafka сама управляет метаданными через встроенный Raft-консенсус, ZooKeeper не нужен.
- Плюсы: проще эксплуатация (одна система), быстрее failover, лучше масштаб по числу партиций.
- В новых версиях KRaft — по умолчанию, ZooKeeper объявлен deprecated.
- Последовательная запись/чтение с диска (append-only) вместо случайного доступа.
- Zero-copy (sendfile): данные идут с диска в сеть, минуя копирование в приложение.
- Page cache ОС вместо своего кэша в JVM-heap.
- Батчинг и сжатие сообщений (
- Партиционирование → горизонтальный параллелизм.
- Broker простой: не трекает статус каждого сообщения, offset хранит consumer.
- Zero-copy (sendfile): данные идут с диска в сеть, минуя копирование в приложение.
- Page cache ОС вместо своего кэша в JVM-heap.
- Батчинг и сжатие сообщений (
linger.ms + compression).- Партиционирование → горизонтальный параллелизм.
- Broker простой: не трекает статус каждого сообщения, offset хранит consumer.
- Партиции = потолок параллелизма: максимум активных consumer'ов в группе = число партиций.
- Оценка: целевой throughput ÷ throughput одной партиции, плюс запас на рост.
- Больше не всегда лучше: много партиций → больше открытых файлов, дольше failover и rebalance, нагрузка на метаданные.
- Число партиций можно увеличить, но не уменьшить; увеличение ломает соответствие key→партиция (порядок).
- Практика: начать с разумного (например, 6–12) и мониторить lag.
- Оценка: целевой throughput ÷ throughput одной партиции, плюс запас на рост.
- Больше не всегда лучше: много партиций → больше открытых файлов, дольше failover и rebalance, нагрузка на метаданные.
- Число партиций можно увеличить, но не уменьшить; увеличение ломает соответствие key→партиция (порядок).
- Практика: начать с разумного (например, 6–12) и мониторить lag.
- Нельзя вечно ретраить «ядовитое» сообщение (poison pill) — застрянет вся партиция.
- Транзиентные ошибки (сеть, таймаут) → retry с backoff.
- Постоянные ошибки (битые данные) → в Dead Letter Queue (отдельный топик) и двигаемся дальше.
- Паттерн: retry-топики с задержкой + DLQ; не коммитить offset, пока сообщение не обработано или не отправлено в DLQ.
- Логировать причину и уметь переиграть из DLQ.
- Транзиентные ошибки (сеть, таймаут) → retry с backoff.
- Постоянные ошибки (битые данные) → в Dead Letter Queue (отдельный топик) и двигаемся дальше.
- Паттерн: retry-топики с задержкой + DLQ; не коммитить offset, пока сообщение не обработано или не отправлено в DLQ.
- Логировать причину и уметь переиграть из DLQ.
- Сервисы общаются через события в Kafka, а не прямыми синхронными вызовами → слабая связанность.
- Producer публикует факт («OrderCreated»), заинтересованные сервисы подписываются и реагируют.
- Плюсы: сервисы не знают друг о друге, легко добавить нового подписчика, буферизация пиков, replay.
- Минусы: eventual consistency, сложнее отлаживать и трейсить, нужно продумывать идемпотентность и порядок.
- Producer публикует факт («OrderCreated»), заинтересованные сервисы подписываются и реагируют.
- Плюсы: сервисы не знают друг о друге, легко добавить нового подписчика, буферизация пиков, replay.
- Минусы: eventual consistency, сложнее отлаживать и трейсить, нужно продумывать идемпотентность и порядок.
- Проблема dual write: записать в БД и отправить в Kafka атомарно нельзя — упадём между ними → рассинхрон.
- Outbox: в той же транзакции БД пишем и бизнес-данные, и запись в таблицу
- Отдельный процесс (или CDC / Debezium) читает
- Гарантия: событие уедет тогда и только тогда, когда закоммичена транзакция БД (at-least-once).
Стандартное решение консистентности БД ↔ Kafka.
- Outbox: в той же транзакции БД пишем и бизнес-данные, и запись в таблицу
outbox. Атомарно.- Отдельный процесс (или CDC / Debezium) читает
outbox и публикует события в Kafka.- Гарантия: событие уедет тогда и только тогда, когда закоммичена транзакция БД (at-least-once).
Стандартное решение консистентности БД ↔ Kafka.
- Ждать глобального порядка по топику (он есть только внутри партиции).
- Коммитить offset до обработки → потеря сообщений при сбое.
- Не делать обработку идемпотентной при at-least-once → дубли.
- Долгая обработка в poll-цикле → превышение
- Игнорировать consumer lag и не иметь DLQ.
- Слишком мало или много партиций; менять их число, забыв про порядок по ключу.
- Возить тяжёлые сообщения вместо ссылки на них (паттерн claim-check).
- Коммитить offset до обработки → потеря сообщений при сбое.
- Не делать обработку идемпотентной при at-least-once → дубли.
- Долгая обработка в poll-цикле → превышение
max.poll.interval.ms → бесконечные rebalance.- Игнорировать consumer lag и не иметь DLQ.
- Слишком мало или много партиций; менять их число, забыв про порядок по ключу.
- Возить тяжёлые сообщения вместо ссылки на них (паттерн claim-check).