Apache Kafka — DE Interview PrepПодготовка к DE-интервью

Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE. Шпаргалка для быстрого повторения перед Data Engineering интервью. Заточена под Senior/Lead DE.
Distributed systemsРаспределённые системы Failure modesРежимы сбоев Streaming at scaleСтриминг на масштабе Anchor: IU Group / Syntea Kafka ingestionЯкорь: IU Group / Syntea Kafka ingestion
🟢 BeginnerНовичок 🔴 AdvancedПродвинутый

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 активный консьюмер на партицию в рамках группы).
Topic "events" (3 partitions, RF=3) P0 ─ off:0 1 2 3 4 ◄─ leader on Broker A (followers B,C) P1 ─ off:0 1 2 3 ◄─ leader on Broker B (followers A,C) P2 ─ off:0 1 2 3 4 5 ◄─ leader on Broker C (followers A,B) ▲ append-only, monotonic offsets, per-partition order
Топик "events" (3 партиции, RF=3) P0 ─ off:0 1 2 3 4 ◄─ лидер на брокере A (реплики B,C) P1 ─ off:0 1 2 3 ◄─ лидер на брокере B (реплики A,C) P2 ─ off:0 1 2 3 4 5 ◄─ лидер на брокере C (реплики A,B) ▲ дозапись в конец, монотонные оффсеты, порядок внутри партиции
Amal anchor On Syntea ingestion I keyed events by student_id so all events for a student landed on one partition (ordered), while spreading load for throughput.
Якорь Amal На ingestion в Syntea я партиционировал события по student_id, чтобы все события одного студента попадали в одну партицию (упорядоченно), при этом распределяя нагрузку ради пропускной способности.
Interviewer will probe "What breaks if you increase partition count?" → keys re-hash to different partitions, so per-key ordering across the resize boundary breaks. Partition count is a careful, near one-way decision.
О чём спросит интервьюер «Что сломается, если увеличить число партиций?» → ключи перехешируются в другие партиции, поэтому порядок по ключу на границе изменения ломается. Число партиций — это аккуратное, почти необратимое решение.

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Два важных таймаута

SettingMeaningIf exceeded
session.timeout.msLiveness via heartbeat threadMember considered dead
max.poll.interval.msMax 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), чтобы быстрый рестарт не запускал ребаланс.
Interviewer will probe "Consumer keeps getting kicked from the group, why?" → per-batch processing exceeds 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() → надёжно записано
  1. Serialize key/value.
  2. Partitioner picks partition (key hash / sticky).
  3. Buffer into a per-partition batch in buffer.memory.
  4. Sender thread flushes when batch.size full or linger.ms elapses.
  5. Broker appends to leader log, replicates to followers.
  6. Ack per acks; idempotence dedupes retries (PID + sequence).
  7. Callback fires (offset on success, or retriable/non-retriable error).
  1. Сериализация ключа/значения.
  2. Партиционер выбирает партицию (хеш ключа / sticky).
  3. Буферизация в батч на партицию в пределах buffer.memory.
  4. Поток sender отправляет, когда заполнен batch.size или истёк linger.ms.
  5. Брокер дописывает в лог лидера, реплицирует на фолловеров.
  6. Ack согласно acks; идемпотентность дедуплицирует ретраи (PID + sequence).
  7. Срабатывает 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/none when 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)Триада надёжности (нужны ВСЕ три)

SettingValueWhy
acksallLeader waits for all ISR before ack
min.insync.replicas2 (with RF=3)Reject writes if <2 in sync (no silent under-replicated writes)
enable.idempotencetrueDedupe retries; bounded retries
НастройкаЗначениеЗачем
acksallЛидер ждёт все ISR перед ack
min.insync.replicas2 (при RF=3)Отклонять запись, если в синхроне <2 (нет тихих недореплицированных записей)
enable.idempotencetrueДедуп ретраев; ограниченные ретраи

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

acksDurabilityRisk
0NoneFire-and-forget, can lose data
1Leader onlyLoss if leader dies pre-replication
allAll ISRHigher latency (tune batching to compensate)
acksНадёжностьРиск
0НетFire-and-forget, можно потерять данные
1Только лидерПотеря, если лидер падает до репликации
allВсе ISRВыше задержка (компенсировать батчингом)
Interviewer will probe "RF=3, 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).
О чём спросит интервьюер «RF=3, min.insync.replicas=3, один брокер упал — что будет?» → ISR (2) < min.isr (3), поэтому продюсеры с acks=all получают ошибку/блокируются. Поэтому min.isr обычно RF−1: переживаем падение одного брокера и продолжаем принимать записи (осознанный выбор в сторону CP).

4. Delivery Semantics4. Семантики доставки

SemanticHowTrade-off
At-most-onceCommit offset before processingNo dupes, possible loss
At-least-once (default, most common)Process then commitNo loss, possible dupes
Exactly-onceIdempotent producer + transactionsNo loss/dupes; cost & complexity
СемантикаКакКомпромисс
At-most-onceКоммит оффсета до обработкиНет дублей, возможна потеря
At-least-once (по умолчанию, самый частый)Сначала обработка, затем коммитНет потерь, возможны дубли
Exactly-onceИдемпотентный продюсер + транзакцииНет потерь/дублей; цена и сложность
Key insight Delivery semantics = where you commit the offset relative to processing. This link gets probed hard at Senior level.
Ключевая мысль Семантика доставки = где ты коммитишь оффсет относительно обработки. На уровне Senior эту связь дожимают сильно.
Amal anchor (honest) Most of my ingestion targeted at-least-once + idempotent/upsert sinks (dedupe at the warehouse on a natural key) because end-to-end EOS adds latency and ops cost that often isn't worth it when the sink can dedupe cheaply.
Якорь Amal (честно) Большая часть моего ingestion была нацелена на at-least-once + идемпотентные/upsert приёмники (дедуп в хранилище по натуральному ключу), потому что сквозной EOS добавляет задержку и операционную стоимость, которые часто не оправданы, если приёмник умеет дешёво дедуплицировать.

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 > 1 AND enable.idempotence=false AND 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 retentionCompaction
policycleanup.policy=deletecleanup.policy=compact
keepsEverything within retention.ms/.bytesAt least latest value per key, forever
use forEvent streams (old events stop mattering)Changelogs / current state of an entity
exampleraw clickstream__consumer_offsets, Streams state, CDC state
Delete retentionКомпакция
политикаcleanup.policy=deletecleanup.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.ms so consumers can observe the delete.
  • Compaction is "latest wins per key," not full event-history preservation.
  • Можно комбинировать: compact,delete.
  • Tombstone = запись с ключом + null-значением → помечает ключ к удалению в компактируемом топике; хранится delete.retention.ms, чтобы консьюмеры успели увидеть удаление.
  • Компакция — это «побеждает последнее значение на ключ», а не сохранение полной истории событий.
Interviewer will probe "Consumer needs every event but topic is compacted?" → don't compact; intermediate values are gone. Use delete retention long enough, or keep a raw event topic + a separate compacted state topic.
О чём спросит интервьюер «Консьюмеру нужны все события, но топик компактируется?» → не компактировать; промежуточные значения уже потеряны. Использовать достаточно длинный delete retention или держать отдельно сырой топик событий + отдельный компактируемый топик состояния.

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 set isolation.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 + sendOffsetsToTransactioncommitTransaction/abortTransaction.
  • The transaction coordinator writes a tx log and inserts commit/abort markers into each partition.
  • read_committed consumers use markers + the Last Stable Offset (LSO) to skip aborted/open-transaction data.
  • Продюсер регистрирует transactional.id → брокер выдаёт PID + epoch (epoch отсекает «зомби» из прошлой инкарнации).
  • beginTransaction → запись + sendOffsetsToTransactioncommitTransaction/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();  // оффсеты + вывод атомарно
}
Critical caveat (say this) Kafka EOS is exactly-once within Kafka. For an external sink (DB/S3/warehouse) you still need an idempotent or transactional sink. Kafka transactions don't extend to arbitrary external systems.
Важная оговорка (произнеси это) Kafka EOS — это exactly-once внутри Kafka. Для внешнего приёмника (БД/S3/хранилище) всё равно нужен идемпотентный или транзакционный приёмник. Транзакции Kafka не распространяются на произвольные внешние системы.
Interviewer will probe "Why not always EOS?" → throughput/latency cost (extra coordination, tx markers, longer commit cycles) and it doesn't cover external sinks. At-least-once + idempotent sink is simpler and usually enough.
О чём спросит интервьюер «Почему не всегда EOS?» → цена по пропускной способности/задержке (доп. координация, tx-маркеры, более долгие циклы коммита), и он не покрывает внешние приёмники. At-least-once + идемпотентный приёмник проще и обычно достаточно.

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=5 safe 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.records vs max.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.

Interviewer will probe "Downside of too many partitions?" → slower recovery/failover, more controller/metadata overhead, more memory for buffers, higher tail latency, __consumer_offsets pressure. Right-size for peak (~2x headroom), not infinity.
О чём спросит интервьюер «Минусы слишком большого числа партиций?» → медленнее восстановление/failover, больше нагрузка на контроллер/метаданные, больше памяти на буферы, выше tail latency, давление на __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 часа ночи в упавшем консьюмере.
  • Безопасно: добавить поле со значением по умолчанию. Ломает: удалить обязательное поле / сменить тип.
Interviewer will probe "Producer adds a required field with no default?" → fails BACKWARD compat for old data/consumers; registry rejects it if compatibility is enforced.
О чём спросит интервьюер «Продюсер добавляет обязательное поле без default?» → ломает BACKWARD-совместимость для старых данных/консьюмеров; registry отклонит, если совместимость включена.
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.
Interviewer will probe "Postgres → warehouse near-real-time, minimal code?" → Debezium source connector + warehouse/S3 sink connector via Connect (Debezium reads the WAL, emits row changes).
О чём спросит интервьюер «Postgres → хранилище почти в реальном времени, минимум кода?» → source-коннектор Debezium + sink-коннектор в хранилище/S3 через Connect (Debezium читает WAL и эмитит изменения строк).

10. Common Failure Modes10. Типичные режимы сбоев

Frame every failure answer as: detect → contain → recover → prevent.

Любой ответ про сбой строй как: обнаружить → локализовать → восстановить → предотвратить.

FailureSymptomHandling
Poison pillDeser/schema error crashes consumer in a loopDLQ with payload + error meta; alert; never block the partition
Consumer lag blowuprecords-lag-max climbingScale consumers ≤ #partitions; fix hot keys; tune fetch; parallelize off poll thread
Hot partition / skewOne partition lags, rest fineFix partition key / composite key / custom partitioner
Rebalance stormGroup churns, processing haltsStatic membership + cooperative rebalancing; reduce per-batch work
Broker/disk/AZ lossUnder-replicated partitionsRF + ISR + rack awareness across AZs
Leader diesBrief unavailabilityNew 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 данные выживают; идемпотентность убирает дубли при ретрае
Amal anchor On Syntea, event volume spiked around exam periods. I sized partitions for peak concurrency and added lag alerting so we caught a slow downstream enrichment step before it became student-visible latency.
Якорь Amal В Syntea объём событий резко рос в экзаменационные периоды. Я заложил число партиций под пиковую конкуренцию и добавил алертинг по лагу, чтобы поймать медленный шаг обогащения ниже по потоку до того, как он превратится в видимую студентам задержку.
Interviewer will probe "How do you reprocess after a bug fix?" → reset offsets on a stopped group (--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 — при условии, что данные ещё в пределах хранения.

12. Flip-to-Reveal Self-Quiz12. Самопроверка: карточки-перевёртыши

Tap a card to flip. Answer out loud first.

Нажми на карточку, чтобы перевернуть. Сначала ответь вслух.

What gives per-entity ordering in Kafka?Что даёт порядок по сущности в Kafka?
tap to revealнажми, чтобы открыть
Use a partition key per entity (e.g. student_id). Order is only guaranteed within a partition, so same key → same partition → ordered.Использовать ключ партиционирования на сущность (например student_id). Порядок гарантирован только внутри партиции, поэтому один ключ → одна партиция → упорядоченно.
The durability triad?Триада надёжности?
tap to revealнажми, чтобы открыть
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.
Delivery semantics depend on what?От чего зависят семантики доставки?
tap to revealнажми, чтобы открыть
Where you commit the offset relative to processing. Before = at-most-once. After = at-least-once. Inside a transaction = exactly-once.Где коммитишь оффсет относительно обработки. До = at-most-once. После = at-least-once. Внутри транзакции = exactly-once.
Idempotent vs transactional producer?Идемпотентный vs транзакционный продюсер?
tap to revealнажми, чтобы открыть
Idempotent = dedupe retries per partition (PID+seq). Transactional = atomic multi-partition writes + offset commit across sessions (consume-process-produce EOS).Идемпотентный = дедуп ретраев на партицию (PID+seq). Транзакционный = атомарная запись в несколько партиций + коммит оффсета через сессии (EOS consume-process-produce).
High watermark?High watermark?
tap to revealнажми, чтобы открыть
Offset all ISR have replicated → the highest offset consumers can read. Gates visibility so you can't read data that might be lost on leader failure.Оффсет, который реплицировали все ISR → максимальный оффсет, доступный консьюмерам для чтения. Ограничивает видимость, чтобы нельзя было прочитать данные, которые могут потеряться при падении лидера.
Compaction vs delete retention?Компакция vs delete retention?
tap to revealнажми, чтобы открыть
Delete = drop old data by time/size (event streams). Compact = keep latest value per key forever (changelogs/state). Tombstone (null value) deletes a key.Delete = удалять старые данные по времени/размеру (потоки событий). Compact = хранить последнее значение на ключ навсегда (changelog/состояние). Tombstone (null-значение) удаляет ключ.
Consumer kicked from group repeatedly — why?Консьюмера регулярно выкидывает из группы — почему?
tap to revealнажми, чтобы открыть
Per-batch processing exceeds 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-потока, увеличить интервал, включить статическое членство.
How can ordering break within a partition?Как порядок ломается внутри партиции?
tap to revealнажми, чтобы открыть
in-flight > 1 + idempotence off + retries → a retried earlier batch lands after a later one. Fix: enable idempotence (order-safe up to 5 in-flight).in-flight > 1 + идемпотентность выключена + ретраи → переотправленный ранний батч приходит после более позднего. Лечение: включить идемпотентность (порядок безопасен до 5 in-flight).
What replaced ZooKeeper?Что заменило ZooKeeper?
tap to revealнажми, чтобы открыть
KRaft (KIP-500): internal Raft controller quorum storing metadata in a Kafka log. Faster failover, more partitions, simpler ops. ZK gone in Kafka 4.x.KRaft (KIP-500): внутренний Raft-кворум контроллеров, хранящий метаданные в Kafka-логе. Быстрее failover, больше партиций, проще эксплуатация. ZK убран в Kafka 4.x.
Poison pill handling?Как обрабатывать poison pill?
tap to revealнажми, чтобы открыть
Route to a dead-letter topic with payload + error metadata, alert, continue. Never block the whole partition on one bad record.Отправить в dead-letter топик с payload + метаданными ошибки, алерт, продолжить. Никогда не блокировать всю партицию из-за одной плохой записи.
Does Kafka EOS cover an external DB sink?Покрывает ли Kafka EOS внешний приёмник-БД?
tap to revealнажми, чтобы открыть
No. EOS is within Kafka. External sinks need their own idempotent upsert or transactional write. Use at-least-once + dedupe on natural key.Нет. EOS — внутри Kafka. Внешним приёмникам нужен свой идемпотентный upsert или транзакционная запись. Использовать at-least-once + дедуп по натуральному ключу.
Downside of too many partitions?Минусы слишком большого числа партиций?
tap to revealнажми, чтобы открыть
Slower recovery/failover, more controller/metadata overhead, more buffer memory, higher tail latency, offsets-topic pressure. Right-size for peak (~2x).Медленнее восстановление/failover, больше нагрузка на контроллер/метаданные, больше памяти на буферы, выше tail latency, давление на топик оффсетов. Подбирать под пик (~2x).