1. Fundamentals: Topics, Partitions, Offsets1. Основы: топики, партиции, оффсеты
- Topic = named logical stream, split into partitions (the unit of parallelism, ordering & replication).
- Partition = append-only, immutable, ordered log on disk.
- Offset = a record's position in a partition. Monotonic, per-partition, never reused.
- Ordering is guaranteed only WITHIN a partition. The partition key decides what's ordered together.
- Routing: explicit partition → key hash (
murmur2(key) % numPartitions) → sticky/round-robin if no key. - Partition count caps consumer parallelism (≤1 active consumer per partition per group).
- Топик (Topic) = именованный логический поток, разбитый на партиции (единица параллелизма, упорядочивания и репликации).
- Партиция (Partition) = неизменяемый лог на диске с дозаписью в конец (append-only) и гарантированным порядком.
- Оффсет (Offset) = позиция записи внутри партиции. Монотонный, отдельный на партицию, не переиспользуется.
- Порядок гарантируется только ВНУТРИ партиции. Ключ партиционирования решает, что упорядочено вместе.
- Маршрутизация: явная партиция → хеш ключа (
murmur2(key) % numPartitions) → sticky/round-robin, если ключа нет. - Число партиций ограничивает параллелизм консьюмеров (≤1 активный консьюмер на партицию в рамках группы).
student_id so all events for a student landed on one partition (ordered), while spreading load for throughput.
student_id, чтобы все события одного студента попадали в одну партицию (упорядоченно), при этом распределяя нагрузку ради пропускной способности.
2. Producers, Consumers, Consumer Groups2. Продюсеры, консьюмеры, группы консьюмеров
Consumer groups & the rebalance protocolГруппы консьюмеров и протокол ребаланса
- Consumer group = consumers sharing a
group.id. The group coordinator assigns each partition to exactly one member. - Offsets committed per group to internal
__consumer_offsets→ that's how a group resumes. - Rebalance = reassign partitions on join/leave/crash, partition-count change, or subscription change.
- Группа консьюмеров = консьюмеры с общим
group.id. Координатор группы назначает каждую партицию ровно одному участнику. - Оффсеты коммитятся на группу во внутренний топик
__consumer_offsets→ так группа возобновляет чтение. - Ребаланс (Rebalance) = переназначение партиций при входе/выходе/падении участника, изменении числа партиций или подписки.
Two timeouts that matterДва важных таймаута
| Setting | Meaning | If exceeded |
|---|---|---|
session.timeout.ms | Liveness via heartbeat thread | Member considered dead |
max.poll.interval.ms | Max gap between poll() calls (progress) | Member kicked → rebalance |
| Настройка | Смысл | При превышении |
|---|---|---|
session.timeout.ms | Признак «жив» через поток heartbeat | Участник считается мёртвым |
max.poll.interval.ms | Макс. интервал между вызовами poll() (прогресс) | Участник исключён → ребаланс |
Rebalance pain at scaleБоль ребаланса на масштабе
- Eager (stop-the-world): everyone revokes all partitions, processing halts → "rebalance storm".
- Fixes: incremental cooperative rebalancing (only moved partitions revoked) + static membership (
group.instance.id) so a quick restart doesn't trigger a rebalance.
- Eager (stop-the-world): все отзывают все партиции, обработка встаёт → «шторм ребалансов».
- Лечится: инкрементальный кооперативный ребаланс (отзываются только перемещаемые партиции) + статическое членство (
group.instance.id), чтобы быстрый рестарт не запускал ребаланс.
max.poll.interval.ms. Fix: lower max.poll.records, move work off the poll thread, or raise the interval.
max.poll.interval.ms. Лечение: уменьшить max.poll.records, вынести работу из poll-потока или увеличить интервал.
Producer internals: send() → durableВнутренности продюсера: send() → надёжно записано
- Serialize key/value.
- Partitioner picks partition (key hash / sticky).
- Buffer into a per-partition batch in
buffer.memory. - Sender thread flushes when
batch.sizefull orlinger.mselapses. - Broker appends to leader log, replicates to followers.
- Ack per
acks; idempotence dedupes retries (PID + sequence). - Callback fires (offset on success, or retriable/non-retriable error).
- Сериализация ключа/значения.
- Партиционер выбирает партицию (хеш ключа / sticky).
- Буферизация в батч на партицию в пределах
buffer.memory. - Поток sender отправляет, когда заполнен
batch.sizeили истёкlinger.ms. - Брокер дописывает в лог лидера, реплицирует на фолловеров.
- Ack согласно
acks; идемпотентность дедуплицирует ретраи (PID + sequence). - Срабатывает callback (оффсет при успехе либо ретраибл/нон-ретраибл ошибка).
send() is async (returns a Future). max.block.ms bounds how long it blocks when buffers are full = backpressure signal.
send() — асинхронный (возвращает Future). max.block.ms ограничивает время блокировки при заполненных буферах = сигнал backpressure.
Offset management: auto vs manual commitУправление оффсетами: авто- vs ручной коммит
- Auto commit (timer-based): simple but commits regardless of processing → loss/dupes on crash. Avoid for anything that matters.
- Manual:
commitSync(safe, blocking) /commitAsync(fast, best-effort). Commit after processing = at-least-once. - Where you commit defines your delivery semantics.
auto.offset.reset=earliest/latest/nonewhen no committed offset exists.
- Авто-коммит (по таймеру): просто, но коммитит независимо от обработки → потери/дубли при падении. Избегать для всего важного.
- Ручной:
commitSync(надёжно, блокирующе) /commitAsync(быстро, best-effort). Коммит после обработки = at-least-once. - Где коммитишь — там и определяются семантики доставки.
auto.offset.reset=earliest/latest/none, когда закоммиченного оффсета нет.
3. Replication, ISR & Durability3. Репликация, ISR и надёжность
- Each partition: 1 leader + (RF−1) followers. Reads/writes go through the leader; followers fetch to stay in sync.
- ISR (In-Sync Replicas) = replicas caught up within
replica.lag.time.max.ms. - High watermark (HW) = offset all ISR have → highest offset a consumer can read.
- Каждая партиция: 1 лидер + (RF−1) фолловеров. Чтение/запись идут через лидера; фолловеры подтягивают данные, чтобы оставаться в синхроне.
- ISR (In-Sync Replicas) = реплики, отставшие не больше чем на
replica.lag.time.max.ms. - High watermark (HW) = оффсет, который есть у всех ISR → максимальный оффсет, доступный консьюмеру для чтения.
The durability triad (you need ALL three)Триада надёжности (нужны ВСЕ три)
| Setting | Value | Why |
|---|---|---|
acks | all | Leader waits for all ISR before ack |
min.insync.replicas | 2 (with RF=3) | Reject writes if <2 in sync (no silent under-replicated writes) |
enable.idempotence | true | Dedupe retries; bounded retries |
| Настройка | Значение | Зачем |
|---|---|---|
acks | all | Лидер ждёт все ISR перед ack |
min.insync.replicas | 2 (при RF=3) | Отклонять запись, если в синхроне <2 (нет тихих недореплицированных записей) |
enable.idempotence | true | Дедуп ретраев; ограниченные ретраи |
Plus unclean.leader.election.enable=false → never elect an out-of-sync replica (which would guarantee loss). Spread replicas across racks/AZs.
Плюс unclean.leader.election.enable=false → никогда не выбирать лидером отставшую реплику (это гарантировало бы потерю данных). Разносить реплики по стойкам/зонам доступности (AZ).
acks comparisonСравнение acks
| acks | Durability | Risk |
|---|---|---|
0 | None | Fire-and-forget, can lose data |
1 | Leader only | Loss if leader dies pre-replication |
all | All ISR | Higher latency (tune batching to compensate) |
| acks | Надёжность | Риск |
|---|---|---|
0 | Нет | Fire-and-forget, можно потерять данные |
1 | Только лидер | Потеря, если лидер падает до репликации |
all | Все ISR | Выше задержка (компенсировать батчингом) |
min.insync.replicas=3, one broker down — what happens?" → ISR (2) < min.isr (3), so acks=all producers error/block. That's why min.isr is usually RF−1: tolerate one broker down and keep accepting writes (deliberate CP-leaning trade-off).
min.insync.replicas=3, один брокер упал — что будет?» → ISR (2) < min.isr (3), поэтому продюсеры с acks=all получают ошибку/блокируются. Поэтому min.isr обычно RF−1: переживаем падение одного брокера и продолжаем принимать записи (осознанный выбор в сторону CP).
4. Delivery Semantics4. Семантики доставки
| Semantic | How | Trade-off |
|---|---|---|
| At-most-once | Commit offset before processing | No dupes, possible loss |
| At-least-once (default, most common) | Process then commit | No loss, possible dupes |
| Exactly-once | Idempotent producer + transactions | No loss/dupes; cost & complexity |
| Семантика | Как | Компромисс |
|---|---|---|
| At-most-once | Коммит оффсета до обработки | Нет дублей, возможна потеря |
| At-least-once (по умолчанию, самый частый) | Сначала обработка, затем коммит | Нет потерь, возможны дубли |
| Exactly-once | Идемпотентный продюсер + транзакции | Нет потерь/дублей; цена и сложность |
5. Ordering Guarantees5. Гарантии порядка
- Guarantee: total order within a partition only. No global/topic-wide order.
- Ordering can break even within a partition if:
max.in.flight.requests.per.connection > 1ANDenable.idempotence=falseAND retries on → a retried earlier batch can land after a later one. - Fix: enable idempotence (preserves order with up to 5 in-flight) or set in-flight to 1.
- Need global order? Single partition (kills parallelism) or carry a logical sequence and reorder downstream.
- Гарантия: полный порядок только внутри партиции. Глобального порядка по топику нет.
- Порядок может сломаться даже внутри партиции, если:
max.in.flight.requests.per.connection > 1Иenable.idempotence=falseИ включены ретраи → переотправленный ранний батч может прийти после более позднего. - Лечение: включить идемпотентность (сохраняет порядок до 5 in-flight) или поставить in-flight = 1.
- Нужен глобальный порядок? Одна партиция (убивает параллелизм) либо нести логический sequence и переупорядочивать на стороне потребителя.
6. Retention & Log Compaction6. Хранение и компакция лога
| Delete retention | Compaction | |
|---|---|---|
| policy | cleanup.policy=delete | cleanup.policy=compact |
| keeps | Everything within retention.ms/.bytes | At least latest value per key, forever |
| use for | Event streams (old events stop mattering) | Changelogs / current state of an entity |
| example | raw clickstream | __consumer_offsets, Streams state, CDC state |
| Delete retention | Компакция | |
|---|---|---|
| политика | cleanup.policy=delete | cleanup.policy=compact |
| хранит | Всё в пределах retention.ms/.bytes | Как минимум последнее значение на ключ, навсегда |
| для чего | Потоки событий (старые события перестают быть важны) | Changelog / текущее состояние сущности |
| пример | сырой clickstream | __consumer_offsets, состояние Streams, состояние CDC |
- Can combine:
compact,delete. - Tombstone = record with key + null value → marks key for deletion in a compacted topic; kept
delete.retention.msso consumers can observe the delete. - Compaction is "latest wins per key," not full event-history preservation.
- Можно комбинировать:
compact,delete. - Tombstone = запись с ключом + null-значением → помечает ключ к удалению в компактируемом топике; хранится
delete.retention.ms, чтобы консьюмеры успели увидеть удаление. - Компакция — это «побеждает последнее значение на ключ», а не сохранение полной истории событий.
7. Exactly-Once: Idempotent Producer & Transactions7. Exactly-Once: идемпотентный продюсер и транзакции
Idempotent vs Transactional producer — the precise differenceИдемпотентный vs транзакционный продюсер — точная разница
- Idempotent (
enable.idempotence=true): broker dedupes producer retries per partition via PID + sequence number. Stops retry dupes within a session. On by default in modern Kafka. Does not give atomic multi-partition writes. - Transactional (
transactional.id): builds on idempotence, adds atomic writes across partitions/topics + offset commit, survives across sessions → true consume-process-produce EOS. Consumers setisolation.level=read_committed.
- Идемпотентный (
enable.idempotence=true): брокер дедуплицирует ретраи продюсера на партицию через PID + sequence number. Убирает дубли от ретраев в рамках сессии. В современной Kafka включён по умолчанию. Не даёт атомарной записи в несколько партиций. - Транзакционный (
transactional.id): строится поверх идемпотентности, добавляет атомарную запись в несколько партиций/топиков + коммит оффсета, переживает рестарты сессий → настоящий EOS по схеме consume-process-produce. Консьюмеры ставятisolation.level=read_committed.
How transactions work on the wireКак транзакции работают «на проводе»
- Producer registers
transactional.id→ broker assigns PID + epoch (epoch fences zombies from a prior incarnation). beginTransaction→ produce +sendOffsetsToTransaction→commitTransaction/abortTransaction.- The transaction coordinator writes a tx log and inserts commit/abort markers into each partition.
read_committedconsumers use markers + the Last Stable Offset (LSO) to skip aborted/open-transaction data.
- Продюсер регистрирует
transactional.id→ брокер выдаёт PID + epoch (epoch отсекает «зомби» из прошлой инкарнации). beginTransaction→ запись +sendOffsetsToTransaction→commitTransaction/abortTransaction.- Координатор транзакций пишет tx-лог и вставляет маркеры commit/abort в каждую партицию.
- Консьюмеры с
read_committedиспользуют маркеры + Last Stable Offset (LSO), чтобы пропускать отменённые/открытые транзакции.
# Java-ish consume-process-produce EOS producer.initTransactions(); while (true) { records = consumer.poll(d); producer.beginTransaction(); for (r : records) producer.send(transform(r)); producer.sendOffsetsToTransaction(offsets, groupMetadata); producer.commitTransaction(); // offsets + output atomic }
# EOS consume-process-produce (псевдо-Java) producer.initTransactions(); while (true) { records = consumer.poll(d); producer.beginTransaction(); for (r : records) producer.send(transform(r)); producer.sendOffsetsToTransaction(offsets, groupMetadata); producer.commitTransaction(); // оффсеты + вывод атомарно }
8. Throughput Tuning8. Тюнинг пропускной способности
ProducerПродюсер
- Batch: raise
batch.size+linger.ms(wait a few ms to fill batches) — small latency cost, big throughput win. - Compress:
compression.type=lz4/zstd(zstd best ratio, lz4 best CPU). max.in.flight.requests.per.connection=5safe with idempotence on.
- Батчинг: поднять
batch.size+linger.ms(подождать пару мс, чтобы наполнить батчи) — небольшая цена по задержке, большой выигрыш в пропускной способности. - Сжатие:
compression.type=lz4/zstd(zstd — лучшая степень сжатия, lz4 — экономнее по CPU). max.in.flight.requests.per.connection=5безопасно при включённой идемпотентности.
ConsumerКонсьюмер
fetch.min.bytes/fetch.max.wait.ms→ bigger fetches.- Balance
max.poll.recordsvsmax.poll.interval.ms; scale consumers to partition count; parallelize off the poll thread.
fetch.min.bytes/fetch.max.wait.ms→ более крупные fetch-запросы.- Балансировать
max.poll.recordsиmax.poll.interval.ms; масштабировать консьюмеров под число партиций; параллелить вне poll-потока.
Why Kafka is fastПочему Kafka быстрая
Sequential disk I/O + OS page cache + zero-copy (sendfile). Give the broker RAM; let the page cache work.
Последовательный дисковый I/O + page cache ОС + zero-copy (sendfile). Дай брокеру RAM; пусть работает page cache.
__consumer_offsets pressure. Right-size for peak (~2x headroom), not infinity.
__consumer_offsets. Подбирать под пик (~2x запас), а не «в бесконечность».
9. Schema Registry & Kafka Connect9. Schema Registry и Kafka Connect
Schema Registry
- Central store for Avro/Protobuf/JSON schemas; message carries a small schema ID, not the whole schema → tiny payloads.
- Enforces compatibility: BACKWARD (default), FORWARD, FULL, NONE. BACKWARD = new schema reads old data.
- Catches breaking changes at produce time, not 2am in a dead consumer.
- Safe: add a field with a default. Breaking: remove a required field / change a type.
- Центральное хранилище схем Avro/Protobuf/JSON; в сообщении лежит маленький schema ID, а не вся схема → крошечный payload.
- Обеспечивает совместимость: BACKWARD (по умолчанию), FORWARD, FULL, NONE. BACKWARD = новая схема читает старые данные.
- Ловит ломающие изменения в момент записи, а не в 2 часа ночи в упавшем консьюмере.
- Безопасно: добавить поле со значением по умолчанию. Ломает: удалить обязательное поле / сменить тип.
Kafka Connect
- Framework for scalable, fault-tolerant integration via reusable connectors: source (DB/CDC → Kafka) and sink (Kafka → warehouse/S3/ES).
- Distributed mode: worker cluster, config-driven, connectors split into tasks for parallelism; offsets/restart handled for you.
- Use Connect for standard integrations (Debezium CDC, S3/JDBC/Elasticsearch sinks) → get retries, scaling, monitoring free.
- Write a custom consumer for non-trivial business logic SMTs can't express.
- Фреймворк для масштабируемой отказоустойчивой интеграции через переиспользуемые коннекторы: source (БД/CDC → Kafka) и sink (Kafka → хранилище/S3/ES).
- Distributed mode: кластер воркеров, управление через конфиги, коннекторы делятся на tasks ради параллелизма; оффсеты/рестарт берёт на себя фреймворк.
- Использовать Connect для стандартных интеграций (Debezium CDC, sink-и в S3/JDBC/Elasticsearch) → бесплатно получаешь ретраи, масштабирование, мониторинг.
- Писать кастомный консьюмер для нетривиальной бизнес-логики, которую не выразить через SMT.
10. Common Failure Modes10. Типичные режимы сбоев
Frame every failure answer as: detect → contain → recover → prevent.
Любой ответ про сбой строй как: обнаружить → локализовать → восстановить → предотвратить.
| Failure | Symptom | Handling |
|---|---|---|
| Poison pill | Deser/schema error crashes consumer in a loop | DLQ with payload + error meta; alert; never block the partition |
| Consumer lag blowup | records-lag-max climbing | Scale consumers ≤ #partitions; fix hot keys; tune fetch; parallelize off poll thread |
| Hot partition / skew | One partition lags, rest fine | Fix partition key / composite key / custom partitioner |
| Rebalance storm | Group churns, processing halts | Static membership + cooperative rebalancing; reduce per-batch work |
| Broker/disk/AZ loss | Under-replicated partitions | RF + ISR + rack awareness across AZs |
| Leader dies | Brief unavailability | New leader from ISR; HW gates visibility; acks=all data survives; idempotence avoids dupes on retry |
| Сбой | Симптом | Обработка |
|---|---|---|
| Poison pill | Ошибка десериализации/схемы крашит консьюмера в цикле | DLQ с payload + метаданными ошибки; алерт; никогда не блокировать партицию |
| Взрыв лага консьюмера | records-lag-max растёт | Масштабировать консьюмеров ≤ числа партиций; чинить горячие ключи; тюнить fetch; параллелить вне poll-потока |
| Горячая партиция / перекос | Одна партиция отстаёт, остальные ок | Чинить ключ партиционирования / составной ключ / кастомный партиционер |
| Шторм ребалансов | Группа «дёргается», обработка встаёт | Статическое членство + кооперативный ребаланс; уменьшить работу на батч |
| Потеря брокера/диска/AZ | Недореплицированные партиции | RF + ISR + распределение по стойкам/AZ |
| Падение лидера | Кратковременная недоступность | Новый лидер из ISR; HW ограничивает видимость; при acks=all данные выживают; идемпотентность убирает дубли при ретрае |
--reset-offsets --to-datetime/--to-earliest) or new group id to replay from retention — requires the data is still within retention.
--reset-offsets --to-datetime/--to-earliest) или взять новый group id, чтобы перечитать из retention — при условии, что данные ещё в пределах хранения.
11. Kafka vs Flink (study gap)11. Kafka vs Flink (пробел для изучения)
- Kafka = durable, partitioned, replayable log / transport (+ light processing via Kafka Streams).
- Flink = true stream-processing engine: event-time, windowing, watermarks, large keyed state with checkpointing, EOS via distributed snapshots (Chandy-Lamport). Kafka is usually Flink's source & sink.
- Kafka = надёжный, партиционированный, перечитываемый лог / транспорт (+ лёгкая обработка через Kafka Streams).
- Flink = полноценный движок потоковой обработки: event-time, окна, watermarks, большое keyed state с чекпоинтами, EOS через распределённые снапшоты (Chandy-Lamport). Kafka обычно служит для Flink источником и приёмником.
12. Flip-to-Reveal Self-Quiz12. Самопроверка: карточки-перевёртыши
Tap a card to flip. Answer out loud first.
Нажми на карточку, чтобы перевернуть. Сначала ответь вслух.
student_id). Order is only guaranteed within a partition, so same key → same partition → ordered.Использовать ключ партиционирования на сущность (например student_id). Порядок гарантирован только внутри партиции, поэтому один ключ → одна партиция → упорядоченно.acks=all + min.insync.replicas=2 (RF=3) + enable.idempotence=true. Plus unclean.leader.election=false.acks=all + min.insync.replicas=2 (RF=3) + enable.idempotence=true. Плюс unclean.leader.election=false.max.poll.interval.ms. Fix: lower max.poll.records, move work off poll thread, raise the interval, use static membership.Обработка батча превышает max.poll.interval.ms. Лечение: уменьшить max.poll.records, вынести работу из poll-потока, увеличить интервал, включить статическое членство.