Перейти к содержанию
Данные и AI

10 вопросов по теме «Data Engineering: Стриминг» на собеседовании

В этом материале — 10 вопросов из русской колоды RecallDeck по теме «Data Engineering: Стриминг». Сначала сформулируйте короткий ответ сами, затем откройте подробный разбор и проверьте примеры, ограничения и отказные случаи.

8 мин чтения10 подробных ответовПроверено 24 августа 2026
Главная мысль

Сначала зафиксируйте grain данных, допущения, метрику и риск leakage, затем обсуждайте модель, инструмент или инфраструктуру.

Вопросы и ответы

10 подробных ответов

01

Объясните базовую модель Kafka: топики, партиции, оффсеты, продюсеры и консьюмеры.

Короткий ответ: Kafka — это распределённый лог коммитов (append-only). Продюсеры дописывают сообщения в конец топика, консьюмеры читают их в своём темпе, а позиция чтения хранится как оффсет. Топик делится на партиции — единицы параллелизма и порядка.

Подробно:

  1. Топик — именованный поток событий; логически это категория сообщений.
  2. Партиция — упорядоченный неизменяемый лог внутри топика. Каждое сообщение получает монотонно растущий оффсет (номер в партиции).
  3. Продюсер пишет в конец партиции; консьюмер читает последовательно и сам двигает свой оффсет (commit), поэтому чтение не удаляет данные.
  4. Брокер — узел кластера, хранящий партиции; данные живут по политике 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

Как партиции и консьюмер-группы дают параллелизм, и как соотносятся число партиций и число консьюмеров?

Короткий ответ: Партиция — единица параллелизма: внутри консьюмер-группы каждую партицию читает ровно один консьюмер. Поэтому эффективный параллелизм ограничен числом партиций — лишние консьюмеры сверх этого простаивают.

Подробно:

  1. Консьюмер-группа — набор консьюмеров с общим group.id, между которыми партиции распределяются (rebalance). Каждая партиция закреплена за одним консьюмером группы.
  2. Порядок гарантируется только внутри партиции. Больше партиций → больше параллелизма, но нет глобального порядка по топику.
  3. Соотношение: консьюмеров ≤ партиций для полной загрузки. Если консьюмеров больше — избыток простаивает без работы.
  4. Масштабирование: партиции задают потолок. Их число легко увеличить, но не уменьшить, и это меняет распределение ключей по партициям.
Партиций Консьюмеров Итог
4 2 по 2 партиции на консьюмера
4 4 по 1 партиции — максимум параллелизма
4 6 4 работают, 2 простаивают

⚠️ Частая ошибка: добавлять консьюмеров ради ускорения, забыв, что потолок — число партиций. Сверх него консьюмеры сидят без назначенных партиций.

03

Чем отличаются at-most-once, at-least-once и exactly-once в Kafka, и как их добиться?

Короткий ответ: Семантику определяет момент коммита оффсета и настройки продюсера. Коммит до обработки даёт at-most-once (можно потерять), коммит после — at-least-once (возможны дубли), а идемпотентный продюсер плюс транзакции дают exactly-once.

Подробно:

  1. At-most-once — оффсет коммитится до обработки. При падении сообщение уже «прочитано», но не обработано → потеря.
  2. At-least-once — оффсет коммитится после успешной обработки. При сбое сообщение перечитается → дубликаты; нужна идемпотентность потребителя.
  3. Exactly-once (EOS)enable.idempotence=true (дедупликация на брокере по producer id + sequence) плюс транзакции (transactional.id), которые атомарно связывают запись результата и коммит оффсета.
Семантика Когда коммитить оффсет Риск
at-most-once до обработки потеря сообщений
at-least-once после обработки дубликаты
exactly-once в транзакции с результатом сложнее, дороже

⚠️ Частая ошибка: коммитить оффсет до фактической обработки ради «скорости» — это молча теряет данные при любом падении консьюмера.

04

Какие гарантии порядка даёт Kafka и как ключ сообщения влияет на маршрутизацию?

Короткий ответ: Порядок гарантируется только внутри одной партиции, не по всему топику. Ключ сообщения определяет партицию: сообщения с одинаковым ключом попадают в одну партицию и читаются по порядку.

Подробно:

  1. Порядок — на уровне партиции. Внутри партиции оффсеты строго возрастают; между партициями порядок не определён.
  2. Маршрутизация по ключу: partition = hash(key) % num_partitions. Один ключ → одна партиция → сохранённый относительный порядок для этого ключа.
  3. Без ключа (key = null) продюсер раскидывает сообщения по партициям (round-robin / sticky) — порядок между ними не гарантирован.
  4. Практика: выбирайте ключом сущность, для которой важен порядок (например 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?

Короткий ответ: Долговечность даёт репликация: каждая партиция копируется на несколько брокеров (replication factor), а ISR — набор реплик, догнавших лидера. Настройка acks балансирует надёжность и латентность, а retention/compaction определяют, как долго живут данные.

Подробно:

  1. Replication factor — сколько копий партиции хранится. RF=3 переживает падение 2 брокеров.
  2. ISR (in-sync replicas) — реплики, синхронные с лидером. acks=all подтверждает запись только когда её приняли все ISR.
  3. Retention — данные хранятся по времени (retention.ms) или размеру; consumer не удаляет их.
  4. Log compaction — альтернатива: для каждого ключа хранится хотя бы последнее значение (снимок состояния).
acks Подтверждение Надёжность / латентность
0 не ждём брокера максимум скорости, потери возможны
1 принял лидер средне; потеря при падении лидера до репликации
all приняли все ISR максимум надёжности, выше латентность

⚠️ Частая ошибка: ставить acks=1 и высокий RF, считая данные защищёнными. Если лидер упадёт до репликации подтверждённой записи, она потеряется — для гарантий нужен acks=all вместе с min.insync.replicas.

06

Какие бывают окна в потоковой обработке (tumbling, sliding, session) и чем event-time отличается от processing-time?

Короткий ответ: Окно группирует события во времени для агрегаций. Tumbling — непересекающиеся фиксированные интервалы, sliding — перекрывающиеся, session — окна по паузам активности. Event-time считает по времени возникновения события, processing-time — по времени обработки.

Подробно:

  1. Tumbling — смежные окна фиксированной длины без перекрытия; каждое событие в ровно одном окне (например, счётчик за каждую минуту).
  2. Sliding — окна фиксированной длины с шагом меньше длины; окна перекрываются, событие попадает в несколько (скользящее среднее).
  3. Session — окно закрывается после паузы бездействия (gap); длина зависит от активности (сессии пользователя).
  4. 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?

Короткий ответ: Водяной знак (watermark) — это оценка «событий раньше момента T мы уже, скорее всего, получили», которая позволяет решить, когда закрывать окно по event-time. Allowed lateness даёт окну дополнительное время принять запоздавшие события после срабатывания watermark.

Подробно:

  1. Проблема: события приходят не по порядку и с задержками (сеть, батчинг, мобильные клиенты). По чистому event-time непонятно, когда результат окна финален.
  2. Watermark движется во времени и говорит движку: «событий с временем < W больше почти не ждём» — по нему окна закрываются и эмитят результат.
  3. Allowed lateness — окно хранит состояние ещё некоторое время после watermark и пересчитывает результат, если приходят поздние события.
  4. За порогом опоздания событие уходит в 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) и чем он лучше батч-поллинга базы?

Короткий ответ: CDC — это захват изменений строк из БД в реальном времени. Debezium читает журнал транзакций (WAL/binlog) и публикует каждую вставку/обновление/удаление как событие в Kafka. В отличие от периодического опроса, CDC даёт низкую задержку, не нагружает базу запросами и не теряет промежуточные изменения.

Подробно:

  1. Как работает: Debezium подключается к логу репликации СУБД (Postgres WAL, MySQL binlog) и превращает коммиты в поток событий — без изменения приложения.
  2. Почему лучше поллинга: опрос по updated_at пропускает промежуточные состояния и удаления, грузит БД тяжёлыми запросами и работает с задержкой в интервал опроса.
  3. Гарантии: события идут в порядке коммитов на партицию, обычно с семантикой at-least-once — потребитель должен быть идемпотентным (upsert по первичному ключу).
  4. Применение: синхронизация в data lake/DWH, инвалидация кэшей, событийная интеграция микросервисов.
┌──────────┐  WAL/binlog  ┌──────────┐  events   ┌────────┐   ┌───────────┐
│ Postgres │ ───────────► │ Debezium │ ────────► │ Kafka  │ ─►│ Lake/DWH  │
└──────────┘              └──────────┘           └────────┘   └───────────┘

⚠️ Частая ошибка: заменять CDC периодическим SELECT ... WHERE updated_at > :last и терять удаления и промежуточные версии строк, а заодно нагружать продакшн-базу.

09

Когда выбирать Kafka, а когда облачную очередь (SQS) или Kinesis/Pulsar?

Короткий ответ: 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.

Источники

Источники и редакционная политика

Материалы RecallDeck сопоставлены с официальной документацией и открытыми публикациями компаний, когда первичный источник доступен. Мы не связаны с упомянутыми работодателями, не публикуем конфиденциальные задания и не продаём места в подборках. Формат найма может меняться — уточняйте его у рекрутера.

От чтения к воспроизведению

Отрепетируйте полный цикл интервью.

RecallDeck возвращает сложные темы по расписанию и помогает удерживать в памяти язык, SQL, архитектуру и поведенческие истории.

Начать подготовку

Продолжить подготовку

Библиотека собеседований RecallDeck

Подробные русские ответы, разборы этапов найма и планы подготовки для российского IT-рынка.

RSS