← Все темы

Kafka

Вопросов: 33

Kafka — распределённый лог событий (distributed event log), он же брокер сообщений для обмена данными между сервисами в реальном времени.

- Модель publish/subscribe: producer'ы пишут события в топики, consumer'ы их читают.
- Данные хранятся на диске и не удаляются после прочтения — их можно перечитывать (replay).
- Горизонтально масштабируется и переживает отказы брокеров.

Зачем на практике: интеграция микросервисов, событийная архитектура, потоковая обработка, сбор логов и метрик, буфер между быстрым producer'ом и медленным consumer'ом.

Три базовые роли в Kafka:

- Producer — пишет (публикует) сообщения в топик. Это твоё приложение.
- Consumer — подписывается на топик и читает сообщения. Тоже твоё приложение.
- Broker — сервер Kafka: хранит данные и обслуживает запросы producer'ов и consumer'ов. Несколько брокеров образуют кластер.

Коротко: producer и consumer — это твой код, broker — сама Kafka.

- Топик (topic) — именованный канал сообщений (например orders). Producer пишет в топик, consumer читает из топика.
- Партиция (partition) — топик физически разбит на N партиций; это единица параллелизма и хранения.
- Внутри одной партиции сообщения строго упорядочены; между партициями порядка нет.
- Больше партиций → выше параллелизм: в группе одну партицию читает один consumer.

Итог: топик — это логика, партиции — это масштаб и параллелизм.

Offset — порядковый номер сообщения внутри партиции (0, 1, 2, …), монотонно растёт.

- Уникален в пределах партиции, а не всего топика.
- Consumer хранит (коммитит) свой offset — «до какого места прочитал», чтобы после перезапуска продолжить, а не начать заново.
- Broker не удаляет сообщение после чтения: несколько групп читают одну партицию независимо, у каждой свой offset.

Сообщение (record) состоит из:

- key — ключ (опционально); определяет партицию и используется в log compaction.
- value — полезная нагрузка (payload), обычно JSON / Avro / Protobuf в виде байтов.
- headers — метаданные (trace-id, тип события и т.п.).
- timestamp — время события или записи.
- Плюс служебное: топик, партиция, offset.

Producer сериализует key/value в байты; Kafka хранит и отдаёт байты и о содержимом ничего не знает.

- Есть ключ → партиция = hash(key) % число_партиций. Одинаковый ключ всегда идёт в одну партицию → сохраняется порядок для этого ключа.
- Нет ключа → сообщения распределяются по партициям равномерно (round-robin / sticky).
- Практика: ключ = сущность, для которой важен порядок (user_id, order_id).
- Осторожно: «горячий» ключ перегружает одну партицию (перекос), а изменение числа партиций ломает соответствие ключ→партиция.

- Каждая партиция — append-only лог: новые сообщения дописываются в конец, старые неизменяемы.
- Лог разбит на сегменты (файлы) плюс индексы offset → позиция для быстрого поиска.
- Последовательная запись и чтение с диска + page cache ОС → очень быстро.
- Данные живут независимо от того, прочитаны они или нет: удаляются по retention, а не после чтения.

Retention — политика, сколько хранить сообщения перед удалением.

- По времени: retention.ms (например, 7 дней).
- По размеру: retention.bytes.
- Что наступит раньше — то и сработает; удаляются целые старые сегменты.
- Настраивается на уровне топика.

Важно: если consumer lag превысит retention — данные удалятся раньше, чем ты их прочитаешь.

- Compaction (cleanup.policy=compact) хранит хотя бы последнее значение для каждого ключа, удаляя старые версии того же ключа.
- Обычный retention удаляет по возрасту/размеру, независимо от ключа.
- Смысл: топик как «снимок текущего состояния» (последний адрес пользователя, актуальный конфиг).
- Удаление ключа — сообщение с этим key и value=null (tombstone).
- Часто используется для changelog- и compacted-топиков.

- У каждой партиции replication.factor копий на разных брокерах.
- Одна реплика — leader: через неё идут все записи и (обычно) чтения. Остальные — followers, копируют данные с лидера.
- Упал брокер с лидером → новым лидером становится один из синхронных followers, данные не теряются.
- Producer и consumer всегда работают с лидером партиции.

ISR — набор реплик (лидер + followers), которые «догнали» лидера и не отстают.

- Новым лидером может стать только реплика из ISR → не теряем подтверждённые данные.
- Отставшая реплика выпадает из ISR; догонит — вернётся.
- Работает в связке с acks=all и min.insync.replicas.

acks — сколько подтверждений от брокеров ждать, прежде чем считать запись успешной. Это компромисс «надёжность ↔ задержка».

- acks=0 — не ждём вовсе. Максимально быстро, но можно потерять (fire-and-forget).
- acks=1 — ждём только лидера. Лидер упал до репликации → сообщение теряется.
- acks=all (-1) — ждём лидера и все реплики из ISR. Максимально надёжно.

Для важных данных: acks=all + идемпотентный producer.

min.insync.replicas — минимум синхронных реплик, при котором разрешена запись с acks=all.

- Реплик в ISR меньше порога → broker отклоняет запись (NotEnoughReplicas), чтобы не принять легко теряемые данные.

Классика durability: replication.factor=3, min.insync.replicas=2, acks=all — переживает отказ одного брокера без потери данных и без остановки записи.

Отказоустойчивость складывается из нескольких механизмов:

- Репликация партиций на разные брокеры.
- Выбор нового лидера из ISR при отказе.
- acks=all + min.insync.replicas → подтверждённое сообщение есть на нескольких репликах.
- Durable append-only лог на диске переживает перезапуск брокера.
- Consumer после сбоя продолжает с закоммиченного offset.

Итог: ни отказ брокера, ни перезапуск consumer'а не теряют подтверждённые данные.

Consumer group — группа consumer'ов с общим group.id, которые вместе читают топик.

- Каждая партиция в группе достаётся ровно одному consumer'у → нагрузка делится.
- Consumer'ов больше, чем партиций → лишние простаивают (параллелизм ограничен числом партиций).
- Разные группы читают один топик независимо, у каждой свои offset (broadcast между группами).

Rebalance — перераспределение партиций между consumer'ами группы.

- Триггеры: consumer вошёл / вышел / упал (перестал слать heartbeat), изменилось число партиций.
- Во время rebalance потребление приостанавливается («stop-the-world») — заметная пауза.
- Частые rebalance — это проблема: долгая обработка превышает max.poll.interval.ms, нестабильные consumer'ы.
- Смягчают: cooperative rebalancing, static membership, корректные таймауты.

- Auto-commit (enable.auto.commit=true) — Kafka периодически коммитит offset сама. Просто, но при сбое возможны потери или дубли (коммит до обработки → потеря).
- Manual commit — коммитишь сам после успешной обработки → at-least-once (дубли возможны, потерь нет).

Практика для надёжности: auto-commit выключить, «обработал → закоммитил», а саму обработку сделать идемпотентной.

Lag = (последний offset в партиции) − (закоммиченный offset 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'а.

- При retry (таймаут, сбой сети) producer может отправить одно сообщение дважды → дубль в партиции.
- Идемпотентный producer (enable.idempotence=true) даёт каждому producer'у PID и порядковые номера; broker отбрасывает повторы.
- Гарантирует «не более одного раза» на запись в партицию и сохраняет порядок при retry.
- Сейчас включён по умолчанию (с acks=all). Дёшево — держать включённым.

- Транзакции — атомарная запись в несколько партиций/топиков плюс коммит offset'ов consumer'а в одной операции. Всё или ничего.
- Основа паттерна consume-process-produce (прочитал → обработал → записал) без дублей = exactly-once.
- Задаётся transactional.id; consumer читает с isolation.level=read_committed (видит только закоммиченное).
- Цена: сложнее и медленнее — включают только там, где реально нужен EOS.

at-least-once означает, что дубли возможны, поэтому обработку делают идемпотентной.

- Уникальный ключ события + проверка «уже обработано» (таблица processed_ids).
- upsert вместо insert; установка значения вместо инкремента (natural idempotency).
- Дедуп по (topic, partition, offset) или по бизнес-ключу.

Часто это проще и надёжнее, чем настраивать полноценный exactly-once через транзакции.

- Порядок гарантирован только внутри одной партиции, а не по топику целиком.
- Нужен порядок для сущности → шли её события с одним ключом (все события 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 хранит просто байты; о структуре договариваются 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 хранила метаданные (брокеры, топики, лидеры) во внешнем ZooKeeper.
- KRaft (KIP-500) — Kafka сама управляет метаданными через встроенный Raft-консенсус, ZooKeeper не нужен.
- Плюсы: проще эксплуатация (одна система), быстрее failover, лучше масштаб по числу партиций.
- В новых версиях KRaft — по умолчанию, ZooKeeper объявлен deprecated.

- Последовательная запись/чтение с диска (append-only) вместо случайного доступа.
- Zero-copy (sendfile): данные идут с диска в сеть, минуя копирование в приложение.
- Page cache ОС вместо своего кэша в JVM-heap.
- Батчинг и сжатие сообщений (linger.ms + compression).
- Партиционирование → горизонтальный параллелизм.
- Broker простой: не трекает статус каждого сообщения, offset хранит consumer.

- Партиции = потолок параллелизма: максимум активных consumer'ов в группе = число партиций.
- Оценка: целевой throughput ÷ throughput одной партиции, плюс запас на рост.
- Больше не всегда лучше: много партиций → больше открытых файлов, дольше failover и rebalance, нагрузка на метаданные.
- Число партиций можно увеличить, но не уменьшить; увеличение ломает соответствие key→партиция (порядок).
- Практика: начать с разумного (например, 6–12) и мониторить lag.

- Нельзя вечно ретраить «ядовитое» сообщение (poison pill) — застрянет вся партиция.
- Транзиентные ошибки (сеть, таймаут) → retry с backoff.
- Постоянные ошибки (битые данные) → в Dead Letter Queue (отдельный топик) и двигаемся дальше.
- Паттерн: retry-топики с задержкой + DLQ; не коммитить offset, пока сообщение не обработано или не отправлено в DLQ.
- Логировать причину и уметь переиграть из DLQ.

- Сервисы общаются через события в Kafka, а не прямыми синхронными вызовами → слабая связанность.
- Producer публикует факт («OrderCreated»), заинтересованные сервисы подписываются и реагируют.
- Плюсы: сервисы не знают друг о друге, легко добавить нового подписчика, буферизация пиков, replay.
- Минусы: eventual consistency, сложнее отлаживать и трейсить, нужно продумывать идемпотентность и порядок.

- Проблема dual write: записать в БД и отправить в Kafka атомарно нельзя — упадём между ними → рассинхрон.
- Outbox: в той же транзакции БД пишем и бизнес-данные, и запись в таблицу outbox. Атомарно.
- Отдельный процесс (или CDC / Debezium) читает outbox и публикует события в Kafka.
- Гарантия: событие уедет тогда и только тогда, когда закоммичена транзакция БД (at-least-once).

Стандартное решение консистентности БД ↔ Kafka.

- Ждать глобального порядка по топику (он есть только внутри партиции).
- Коммитить offset до обработки → потеря сообщений при сбое.
- Не делать обработку идемпотентной при at-least-once → дубли.
- Долгая обработка в poll-цикле → превышение max.poll.interval.ms → бесконечные rebalance.
- Игнорировать consumer lag и не иметь DLQ.
- Слишком мало или много партиций; менять их число, забыв про порядок по ключу.
- Возить тяжёлые сообщения вместо ссылки на них (паттерн claim-check).