Сначала зафиксируйте grain данных, допущения, метрику и риск leakage, затем обсуждайте модель, инструмент или инфраструктуру.
Вопросы и ответы
10 подробных ответов
01Объясните базовую модель Kafka: топики, партиции, оффсеты, продюсеры и консьюмеры.
middle
Короткий ответ: Kafka — это распределённый лог коммитов (append-only). Продюсеры дописывают сообщения в конец топика, консьюмеры читают их в своём темпе, а позиция чтения хранится как оффсет. Топик делится на партиции — единицы параллелизма и порядка.
Подробно:
- Топик — именованный поток событий; логически это категория сообщений.
- Партиция — упорядоченный неизменяемый лог внутри топика. Каждое сообщение получает монотонно растущий оффсет (номер в партиции).
- Продюсер пишет в конец партиции; консьюмер читает последовательно и сам двигает свой оффсет (commit), поэтому чтение не удаляет данные.
- Брокер — узел кластера, хранящий партиции; данные живут по политике retention, а не «до первого прочтения».
Топик "orders"
Partition 0: [0][1][2][3][4]───► новые пишутся в конец
Partition 1: [0][1][2]
Partition 2: [0][1][2][3]
▲
оффсет консьюмера = 2 (следующим прочитает 3)
⚠️ Частая ошибка: думать, что консьюмер «забирает» сообщение и оно исчезает, как в очереди. В Kafka лог остаётся, много консьюмер-групп читают его независимо, а удаляет данные только retention.
02Как партиции и консьюмер-группы дают параллелизм, и как соотносятся число партиций и число консьюмеров?
senior
Короткий ответ: Партиция — единица параллелизма: внутри консьюмер-группы каждую партицию читает ровно один консьюмер. Поэтому эффективный параллелизм ограничен числом партиций — лишние консьюмеры сверх этого простаивают.
Подробно:
- Консьюмер-группа — набор консьюмеров с общим
group.id, между которыми партиции распределяются (rebalance). Каждая партиция закреплена за одним консьюмером группы. - Порядок гарантируется только внутри партиции. Больше партиций → больше параллелизма, но нет глобального порядка по топику.
- Соотношение: консьюмеров ≤ партиций для полной загрузки. Если консьюмеров больше — избыток простаивает без работы.
- Масштабирование: партиции задают потолок. Их число легко увеличить, но не уменьшить, и это меняет распределение ключей по партициям.
| Партиций | Консьюмеров | Итог |
|---|---|---|
| 4 | 2 | по 2 партиции на консьюмера |
| 4 | 4 | по 1 партиции — максимум параллелизма |
| 4 | 6 | 4 работают, 2 простаивают |
⚠️ Частая ошибка: добавлять консьюмеров ради ускорения, забыв, что потолок — число партиций. Сверх него консьюмеры сидят без назначенных партиций.
03Чем отличаются at-most-once, at-least-once и exactly-once в Kafka, и как их добиться?
senior
Короткий ответ: Семантику определяет момент коммита оффсета и настройки продюсера. Коммит до обработки даёт at-most-once (можно потерять), коммит после — at-least-once (возможны дубли), а идемпотентный продюсер плюс транзакции дают exactly-once.
Подробно:
- At-most-once — оффсет коммитится до обработки. При падении сообщение уже «прочитано», но не обработано → потеря.
- At-least-once — оффсет коммитится после успешной обработки. При сбое сообщение перечитается → дубликаты; нужна идемпотентность потребителя.
- Exactly-once (EOS) —
enable.idempotence=true(дедупликация на брокере по producer id + sequence) плюс транзакции (transactional.id), которые атомарно связывают запись результата и коммит оффсета.
| Семантика | Когда коммитить оффсет | Риск |
|---|---|---|
| at-most-once | до обработки | потеря сообщений |
| at-least-once | после обработки | дубликаты |
| exactly-once | в транзакции с результатом | сложнее, дороже |
⚠️ Частая ошибка: коммитить оффсет до фактической обработки ради «скорости» — это молча теряет данные при любом падении консьюмера.
04Какие гарантии порядка даёт Kafka и как ключ сообщения влияет на маршрутизацию?
middle
Короткий ответ: Порядок гарантируется только внутри одной партиции, не по всему топику. Ключ сообщения определяет партицию: сообщения с одинаковым ключом попадают в одну партицию и читаются по порядку.
Подробно:
- Порядок — на уровне партиции. Внутри партиции оффсеты строго возрастают; между партициями порядок не определён.
- Маршрутизация по ключу:
partition = hash(key) % num_partitions. Один ключ → одна партиция → сохранённый относительный порядок для этого ключа. - Без ключа (
key = null) продюсер раскидывает сообщения по партициям (round-robin / sticky) — порядок между ними не гарантирован. - Практика: выбирайте ключом сущность, для которой важен порядок (например
user_id,order_id), чтобы её события шли последовательно.
key="user-42" ─hash─► Partition 1 [событие A][событие B][событие C] (порядок сохранён)
key="user-99" ─hash─► Partition 0
key=null ─round-robin─► любая партиция
⚠️ Частая ошибка: ждать глобального порядка по всему топику. С несколькими партициями его нет — порядок есть только в пределах ключа/партиции.
05Как Kafka обеспечивает надёжность: replication factor, ISR, acks и retention/compaction?
middle
Короткий ответ: Долговечность даёт репликация: каждая партиция копируется на несколько брокеров (replication factor), а ISR — набор реплик, догнавших лидера. Настройка acks балансирует надёжность и латентность, а retention/compaction определяют, как долго живут данные.
Подробно:
- Replication factor — сколько копий партиции хранится. RF=3 переживает падение 2 брокеров.
- ISR (in-sync replicas) — реплики, синхронные с лидером.
acks=allподтверждает запись только когда её приняли все ISR. - Retention — данные хранятся по времени (
retention.ms) или размеру; consumer не удаляет их. - Log compaction — альтернатива: для каждого ключа хранится хотя бы последнее значение (снимок состояния).
| acks | Подтверждение | Надёжность / латентность |
|---|---|---|
| 0 | не ждём брокера | максимум скорости, потери возможны |
| 1 | принял лидер | средне; потеря при падении лидера до репликации |
| all | приняли все ISR | максимум надёжности, выше латентность |
⚠️ Частая ошибка: ставить acks=1 и высокий RF, считая данные защищёнными. Если лидер упадёт до репликации подтверждённой записи, она потеряется — для гарантий нужен acks=all вместе с min.insync.replicas.
06Какие бывают окна в потоковой обработке (tumbling, sliding, session) и чем event-time отличается от processing-time?
senior
Короткий ответ: Окно группирует события во времени для агрегаций. Tumbling — непересекающиеся фиксированные интервалы, sliding — перекрывающиеся, session — окна по паузам активности. Event-time считает по времени возникновения события, processing-time — по времени обработки.
Подробно:
- Tumbling — смежные окна фиксированной длины без перекрытия; каждое событие в ровно одном окне (например, счётчик за каждую минуту).
- Sliding — окна фиксированной длины с шагом меньше длины; окна перекрываются, событие попадает в несколько (скользящее среднее).
- Session — окно закрывается после паузы бездействия (gap); длина зависит от активности (сессии пользователя).
- Event-time vs processing-time: event-time даёт корректные результаты при задержках и переупорядочивании, но требует водяных знаков; processing-time проще и быстрее, но искажает результат при лагах.
Ось времени ──────────────────────────►
Tumbling: [ 0–1м ][ 1–2м ][ 2–3м ]
Sliding: [ 0–1м ]
[ 30с–1:30 ]
[ 1–2м ]
Session: [ai a2 a3] (gap) [b1 b2]
⚠️ Частая ошибка: агрегировать по processing-time и удивляться «неверным» суточным итогам, когда события пришли с задержкой. Для точных бизнес-метрик считайте по event-time.
07Как обрабатывать поздние и переупорядоченные события: что такое водяные знаки и allowed lateness?
senior
Короткий ответ: Водяной знак (watermark) — это оценка «событий раньше момента T мы уже, скорее всего, получили», которая позволяет решить, когда закрывать окно по event-time. Allowed lateness даёт окну дополнительное время принять запоздавшие события после срабатывания watermark.
Подробно:
- Проблема: события приходят не по порядку и с задержками (сеть, батчинг, мобильные клиенты). По чистому event-time непонятно, когда результат окна финален.
- Watermark движется во времени и говорит движку: «событий с временем < W больше почти не ждём» — по нему окна закрываются и эмитят результат.
- Allowed lateness — окно хранит состояние ещё некоторое время после watermark и пересчитывает результат, если приходят поздние события.
- За порогом опоздания событие уходит в side output (dead-letter) или отбрасывается — компромисс между точностью и латентностью/памятью.
события по event-time: e(10:00) e(10:02) e(09:59←поздно) e(10:03)
watermark ───────────────────────► 10:02
└ окно 10:00–10:01 закрыто; e(09:59) в пределах
allowed lateness → окно пересчитано
⚠️ Частая ошибка: ставить слишком «строгий» watermark ради низкой латентности и молча терять поздние события — либо, наоборот, огромный allowed lateness, раздувающий состояние в памяти.
08Что такое Change Data Capture с Kafka (например, Debezium) и чем он лучше батч-поллинга базы?
middle
Короткий ответ: CDC — это захват изменений строк из БД в реальном времени. Debezium читает журнал транзакций (WAL/binlog) и публикует каждую вставку/обновление/удаление как событие в Kafka. В отличие от периодического опроса, CDC даёт низкую задержку, не нагружает базу запросами и не теряет промежуточные изменения.
Подробно:
- Как работает: Debezium подключается к логу репликации СУБД (Postgres WAL, MySQL binlog) и превращает коммиты в поток событий — без изменения приложения.
- Почему лучше поллинга: опрос по
updated_atпропускает промежуточные состояния и удаления, грузит БД тяжёлыми запросами и работает с задержкой в интервал опроса. - Гарантии: события идут в порядке коммитов на партицию, обычно с семантикой at-least-once — потребитель должен быть идемпотентным (upsert по первичному ключу).
- Применение: синхронизация в data lake/DWH, инвалидация кэшей, событийная интеграция микросервисов.
┌──────────┐ WAL/binlog ┌──────────┐ events ┌────────┐ ┌───────────┐
│ Postgres │ ───────────► │ Debezium │ ────────► │ Kafka │ ─►│ Lake/DWH │
└──────────┘ └──────────┘ └────────┘ └───────────┘
⚠️ Частая ошибка: заменять CDC периодическим SELECT ... WHERE updated_at > :last и терять удаления и промежуточные версии строк, а заодно нагружать продакшн-базу.
09Когда выбирать Kafka, а когда облачную очередь (SQS) или Kinesis/Pulsar?
middle
Короткий ответ: Kafka — для высокопропускного, переигрываемого лога событий с несколькими независимыми консьюмерами. SQS — простая управляемая очередь для распределения задач без реплея. Kinesis — «Kafka-подобный» managed-стриминг в AWS, Pulsar — альтернатива с мультиарендностью и разделением compute/storage.
Подробно:
- Kafka — большой throughput, хранение и повторное чтение лога, много консьюмер-групп, экосистема (Connect, Streams). Ценой операционной сложности (или managed вроде Confluent/MSK).
- SQS — очередь задач: сообщение прочитал-удалил, нет реплея и порядка (кроме FIFO), зато почти нулевые эксплуатационные затраты.
- Kinesis / Pulsar — стриминг с ретеншеном как у Kafka; Kinesis нативен в AWS, Pulsar даёт мультиарендность и независимое масштабирование хранилища.
| Критерий | Kafka | SQS | Kinesis |
|---|---|---|---|
| Модель | лог событий | очередь задач | лог событий (managed) |
| Повторное чтение | да (retention) | нет | да (retention) |
| Много консьюмеров | да, группы | конкурирующие воркеры | да, shards |
| Эксплуатация | сложнее | минимальная | managed AWS |
⚠️ Частая ошибка: тащить Kafka туда, где нужна простая очередь задач без реплея — получаете лишнюю операционную нагрузку там, где хватило бы SQS.
10В чём ключевое отличие Flink от Spark Structured Streaming (микробатч против настоящего стриминга)?
senior
Короткий ответ: Flink — настоящий потоковый движок: обрабатывает каждое событие по мере поступления, что даёт латентность в миллисекундах. Spark Structured Streaming исторически работает микробатчами — копит события за короткий интервал и обрабатывает пачкой, что проще, но добавляет задержку.
Подробно:
- Модель обработки: Flink — event-at-a-time (истинный стриминг); Spark — микробатчи (мини-джобы каждые N мс/сек). У Spark есть режим Continuous Processing, но основной путь — микробатч.
- Латентность: Flink — суб-секундная/миллисекунды; микробатч Spark — в пределах интервала батча.
- Состояние и время: у Flink развитое управление состоянием, event-time, watermarks и точные окна — это его сильная сторона.
- Экосистема: Spark удобен, если у команды уже есть Spark для батча (единый стек batch + streaming); Flink выбирают, когда критична низкая задержка и сложная событийная логика.
| Аспект | Flink | Spark Structured Streaming |
|---|---|---|
| Модель | по событию (true streaming) | микробатч |
| Латентность | миллисекунды | интервал батча |
| Состояние/окна | богатые, event-time | есть, но исторически проще |
| Когда брать | low-latency стриминг | уже есть Spark-стек |
⚠️ Частая ошибка: называть Spark Structured Streaming «настоящим» стримингом наравне с Flink. По умолчанию это микробатч, и латентность ограничена интервалом батча.
Источники
Источники и редакционная политика
Материалы RecallDeck сопоставлены с официальной документацией и открытыми публикациями компаний, когда первичный источник доступен. Мы не связаны с упомянутыми работодателями, не публикуем конфиденциальные задания и не продаём места в подборках. Формат найма может меняться — уточняйте его у рекрутера.
От чтения к воспроизведению
Отрепетируйте полный цикл интервью.
RecallDeck возвращает сложные темы по расписанию и помогает удерживать в памяти язык, SQL, архитектуру и поведенческие истории.