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

Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE. Шпаргалка для быстрого повторения перед Data Engineering интервью. Заточена под Senior/Lead DE.
⚠ STUDY-GAP: not in recent stack⚠ ПРОБЕЛ: не в недавнем стеке Bridge from: Kafka · Scala/Spark · AirflowМост от: Kafka · Scala/Spark · Airflow Senior/Lead depthУровень Senior/Lead

🎯 Front-Loaded Must-Knows🎯 Главное, что надо знать

The 8 things you must say cleanly8 вещей, которые нужно сказать чётко

  • Event time vs processing time: logic keys off the timestamp in the event so replays are deterministic. Processing time = wall-clock = fast but non-deterministic.
  • Watermark = "no earlier events expected past this point." Triggers event-time windows; is the min across inputs.
  • Windows: tumbling (fixed, no overlap), sliding (overlap), session (gap-based), global (custom trigger).
  • Exactly-once = state consistency via aligned checkpoint barriers (Chandy-Lamport). End-to-end needs a transactional/idempotent sink.
  • State backends: HashMap (heap, small/fast) vs EmbeddedRocksDB (disk, huge state, incremental checkpoints).
  • Checkpoint (auto, recovery) vs Savepoint (manual, portable, for upgrades/rescale).
  • Keyed state after keyBy; rescaled via key groups; bound it with TTL.
  • Canonical pipeline: Kafka → Flink (exactly-once stateful) → Iceberg → Trino/Spark.
  • Event time vs processing time: логика опирается на временную метку в событии, поэтому повторы детерминированы. Processing time = wall-clock = быстро, но не детерминировано.
  • Watermark = «событий с меткой времени раньше этого момента больше не ожидается». Запускает окна event-time; равен min по всем входам.
  • Окна (Windows): tumbling (фиксированные, без перекрытия), sliding (с перекрытием), session (на основе паузы), global (кастомный триггер).
  • Exactly-once = согласованность состояния через выровненные checkpoint barriers (Chandy-Lamport). Сквозная требует транзакционного/идемпотентного приёмника.
  • State backends: HashMap (heap, маленькое/быстрое) vs EmbeddedRocksDB (диск, огромное состояние, инкрементальные checkpoints).
  • Checkpoint (авто, восстановление) vs Savepoint (ручной, портируемый, для апгрейдов/масштабирования).
  • Keyed state после keyBy; масштабируется через key groups; ограничь через TTL.
  • Канонический пайплайн: Kafka → Flink (exactly-once stateful) → Iceberg → Trino/Spark.
YOUR ANCHOR LINEТВОЯ ЯКОРНАЯ ФРАЗА

"I haven't run Flink in production, but I built Kafka ingestion at IU Group, owned a Scala/Spark stateful training pipeline on Hadoop at Front Tier, and ran Airflow migrations at McMakler. Event time, watermarks, state and exactly-once are the same failure-mode reasoning I already did on those systems, so let me reason from first principles."

«Я не запускал Flink в продакшене, но строил ingestion на Kafka в IU Group, владел stateful-пайплайном на Scala/Spark поверх Hadoop во Front Tier и управлял миграциями Airflow в McMakler. Event time, watermarks, состояние и exactly-once — это те же рассуждения о режимах сбоев, которые я уже делал на тех системах, поэтому дай мне порассуждать с первых принципов.»

GAP STRATEGYСТРАТЕГИЯ ПРОБЕЛА

Be honest about no prod Flink. Convert it: offset/exactly-once reasoning (Kafka) + stateful Scala/Spark + orchestration. Then show a concrete de-risk plan: UIDs on every stateful op, state TTL, checkpoint monitoring, MiniCluster integration tests, shadow on-call.

Честно признай отсутствие Flink в проде. Конвертируй: рассуждения про offset/exactly-once (Kafka) + stateful Scala/Spark + оркестрация. Затем покажи конкретный план снижения рисков: UIDs на каждой stateful-операции, state TTL, мониторинг checkpoint, MiniCluster integration tests, теневой on-call.

⏱ Stream vs Batch & Event Time⏱ Поток vs батч и Event Time

Stream vs BatchПоток vs батч

  • Batch = bounded dataset, runs to completion, latency in min/hours.
  • Stream = unbounded dataset, runs forever, latency in ms.
  • Flink treats batch as a special case of streaming (bounded stream). Unified runtime.
  • Streaming needs answers to: out-of-order data, when to emit, fault tolerance with running state.
  • Batch = ограниченный датасет, отрабатывает до конца, задержка в мин/часах.
  • Stream = неограниченный датасет, работает вечно, задержка в мс.
  • Flink трактует батч как частный случай стриминга (ограниченный поток). Единая среда выполнения.
  • Стриминг требует ответов на: данные вне порядка, когда выдавать результат, отказоустойчивость при наличии состояния.

Three notions of timeТри понятия времени

TypeMeaningTrade-off
Eventtimestamp in the recordcorrect + replayable; needs watermarks → latency
Ingestionassigned at source opmiddle ground
Processingoperator wall-clockfastest; non-deterministic, wrong under skew
ТипСмыслКомпромисс
Eventметка времени в записикорректно + повторяемо; требует watermarks → задержка
Ingestionназначается в source-операторезолотая середина
Processingwall-clock операторабыстрее всего; не детерминировано, неверно при перекосе
INTERVIEWER WILL PROBEО ЧЁМ СПРОСИТ ИНТЕРВЬЮЕР

"Why event time if it adds latency?" → Determinism. Backfills/replays produce identical results because logic uses the embedded timestamp, not when the job happened to run.

«Зачем event time, если добавляет задержку?» → Детерминизм. Backfill-ы/повторы дают идентичные результаты, потому что логика использует встроенную метку времени, а не момент запуска джоба.

ANCHORЯКОРЬ

On the bank default-prediction pipeline, features had to be computed as-of the loan-decision time, not as-of run time. Event time is that as-of discipline made first-class.

В пайплайне прогнозирования дефолтов для банка фичи должны были вычисляться на момент принятия решения о займе, а не на момент запуска. Event time — это та самая дисциплина «на момент», сделанная first-class.

🌊 Watermarks🌊 Watermarks

What a watermark actually isЧто такое watermark на самом деле

A marker flowing in the stream: W(t) = "no event with timestamp ≤ t will arrive after this." It is Flink's notion of event-time progress and is what fires event-time windows (every window whose end ≤ W).

Маркер, текущий в потоке: W(t) = «событий с меткой времени ≤ t после этого момента не будет». Это понятие прогресса event-time в Flink и то, что запускает event-time окна (каждое окно, конец которого ≤ W).

// bounded out-of-orderness: watermark = maxSeenTs - 5s
WatermarkStrategy
  .forBoundedOutOfOrderness(Duration.ofSeconds(5))
  .withTimestampAssigner((event, ts) -> event.getEventTime())
  .withIdleness(Duration.ofMinutes(1)); // don't let idle partition stall WM
  • Per-partition watermarks: operator watermark = minimum across all inputs/partitions. Slowest input gates progress.
  • Idle source stalls the min forever → withIdleness() marks it idle.
  • The knob: bigger out-of-orderness = more correctness for late data, more latency. Smaller = faster, more drops.
  • Watermarks по партициям: watermark оператора = минимум по всем входам/партициям. Самый медленный вход блокирует прогресс.
  • Idle source блокирует минимум навсегда → withIdleness() помечает его idle.
  • Ручка настройки: больше out-of-orderness = больше корректности для поздних данных, больше задержка. Меньше = быстрее, больше сброшенных данных.
event-time axis → events: 1 4 2 7 3(late) 9 ▲ arrives after WM=5 → "late" watermark advances: ... W=5 ........... W=9 └ fires all windows with end ≤ 5
ось event-time → события: 1 4 2 7 3(опоздало) 9 ▲ пришло после WM=5 → «опоздало» watermark двигается: ... W=5 ........... W=9 └ запускает все окна с концом ≤ 5
PROBEВОПРОС

"One Kafka partition is silent — what happens to your windows?" → its watermark never advances, min stays low, windows never fire. Fix: withIdleness.

«Одна партиция Kafka молчит — что будет с твоими окнами?» → её watermark не двигается, минимум остаётся низким, окна не запускаются. Лечение: withIdleness.

🪟 Windows🪟 Окна

WindowShapeUse for
Tumblingfixed size, no overlap, 1 event → 1 windowperiodic buckets ("count per 1 min")
Slidingsize + slide; overlap; event in many windowsmoving averages ("5-min count every 1 min")
Sessiondynamic, closes after gap of inactivityuser sessionization ("clicks until 30-min idle")
Globalone window per key, custom Trigger fires itbuilding block / custom logic
ОкноФормаДля чего
Tumblingфиксированный размер, без перекрытия, 1 событие → 1 окнопериодические бакеты («счёт за 1 мин»)
Slidingразмер + слайд; перекрытие; событие во многих окнахскользящие средние («счёт за 5 мин каждую 1 мин»)
Sessionдинамическое, закрывается после паузысессионизация пользователей («клики до 30-мин простоя»)
Globalодно окно на ключ, кастомный Trigger запускает егобазовый блок / кастомная логика

Machinery: Assigner (which window) · Trigger (when to fire) · Evictor (optional removal) · allowed lateness.

Механизм: Assigner (какое окно) · Trigger (когда запускать) · Evictor (опциональное удаление) · allowed lateness.

stream.keyBy(e -> e.userId())
  .window(TumblingEventTimeWindows.of(Time.minutes(1)))
  .allowedLateness(Time.seconds(30))
  .sideOutputLateData(lateTag)   // don't drop late data silently
  .aggregate(new CountAgg());

Late data handlingОбработка опоздавших данных

  • Past the watermark = late. Default: dropped.
  • Allowed lateness: keep window state open longer; late events re-fire / update result.
  • Side output (OutputTag): route too-late events to a separate stream instead of losing them.
  • Позже watermark = опоздало. По умолчанию: сброшено.
  • Allowed lateness: держать состояние окна открытым дольше; опоздавшие события перезапускают / обновляют результат.
  • Side output (OutputTag): направить слишком поздние события в отдельный поток вместо потери.
ANCHORЯКОРЬ

Sessionization is the streaming-native version of the gaps-and-islands SQL I wrote for user behavior at McMakler.

Сессионизация — это streaming-native версия SQL-паттерна gaps-and-islands, который я писал для поведения пользователей в McMakler.

🗃 State, Keyed State & Backends🗃 Состояние, Keyed State и бэкенды

Keyed vs Operator stateKeyed vs Operator state

  • Keyed state (after keyBy): one independent instance per key. Types: ValueState, ListState, MapState, ReducingState, AggregatingState.
  • Operator state: scoped to an operator instance, not a key (e.g. Kafka source offsets). Redistributed on rescale.
  • Keyed state (после keyBy): один независимый экземпляр на ключ. Типы: ValueState, ListState, MapState, ReducingState, AggregatingState.
  • Operator state: привязан к экземпляру оператора, а не к ключу (например, оффсеты Kafka source). Перераспределяется при масштабировании.

State backendsState backends

BackendWhereWhen
HashMapJVM heapsmall/fast state, low latency
EmbeddedRocksDBlocal disk, off-heap, serializedhuge state > memory; incremental checkpoints
BackendГдеКогда
HashMapJVM heapмаленькое/быстрое состояние, низкая задержка
EmbeddedRocksDBлокальный диск, off-heap, сериализованноеогромное состояние > памяти; инкрементальные checkpoints

Checkpoint storage (HDFS/S3) is separate: it's where the snapshot lands.

Checkpoint storage (HDFS/S3) — отдельно: это куда попадает снапшот.

Rescaling & TTLМасштабирование и TTL

  • Key groups = atomic unit of keyed-state redistribution. Change parallelism → key groups reassigned across subtasks.
  • Max parallelism fixes the number of key groups → caps how far you can rescale without a state migration. Set it deliberately at job creation.
  • State TTL (StateTtlConfig): expire old keys or state grows forever → slow recovery, disk blowup. Essential at per-user, billions-of-keys scale.
  • Key groups = атомарная единица перераспределения keyed-state. Изменение parallelism → key groups переназначаются между subtasks.
  • Max parallelism фиксирует число key groups → ограничивает, насколько можно масштабироваться без миграции состояния. Устанавливай осознанно при создании джоба.
  • State TTL (StateTtlConfig): устаревание старых ключей, иначе состояние растёт вечно → медленное восстановление, взрыв диска. Критично на масштабе миллиардов ключей по пользователям.
PROBEВОПРОС

"What if state grows unbounded?" → State TTL + RocksDB incremental checkpoints; bound joins with intervals.

«Что, если состояние растёт неограниченно?» → State TTL + RocksDB инкрементальные checkpoints; ограничь join-ы интервалами.

ANCHORЯКОРЬ

Large keyed state in RocksDB is a partitioned, spill-to-disk keyed store. Entity resolution across 6.2B records at Meta drilled the same instinct: key the data, keep per-key state bounded, deterministic recovery.

Большое keyed state в RocksDB — это партиционированное, spill-to-disk хранилище по ключам. Entity resolution по 6.2B записям в Meta тренировал тот же инстинкт: ключить данные, держать состояние на ключ ограниченным, детерминированное восстановление.

💾 Checkpoints & Savepoints💾 Checkpoints и Savepoints

CheckpointSavepoint
Purposeautomatic fault tolerancemanual operational snapshot
Triggerperiodic, by Flinkon-demand, by you/CI
OwnershipFlink-managed, may be cleaned upuser-owned, retained, portable
Use caserecover from crashupgrade code, rescale, migrate cluster/version
CheckpointSavepoint
Цельавтоматическая отказоустойчивостьручной операционный снапшот
Триггерпериодически, Flink-омпо требованию, тобой/CI
Владениеуправляется Flink, может быть очищенв собственности пользователя, сохранён, портируемый
Use caseвосстановление после паденияапгрейд кода, масштабирование, миграция кластера/версии

Workflow: stop job with savepoint → deploy new code → resume from savepoint. State survives a deploy as long as schema + UIDs are compatible.

Workflow: остановить джоб с savepoint → задеплоить новый код → продолжить из savepoint. Состояние переживает деплой, пока схема + UID совместимы.

PROBE — classic footgunВОПРОС — классическая ловушка

"What lets state survive a code change?" → stable operator UIDs .uid("..."). On restore Flink maps saved state to operators by ID. Auto IDs change when topology changes → restore fails / state silently dropped. Set .uid() on every stateful op from day one.

«Что позволяет состоянию пережить изменение кода?» → стабильные UIDs операторов .uid("..."). При восстановлении Flink сопоставляет сохранённое состояние с операторами по ID. Авто-ID меняются при изменении топологии → восстановление падает / состояние тихо теряется. Ставь .uid() на каждую stateful-операцию с первого дня.

ANCHORЯКОРЬ

Monitoring checkpoint duration / size / alignment / failures is SLO thinking — same as the DQ monitors and completeness alerts I built on Meta pipelines: instrument the thing that silently degrades before it pages you.

Мониторинг длительности / размера / alignment / сбоев checkpoint — это SLO-мышление, как DQ-мониторы и алерты на полноту, которые я построил на пайплайнах в Meta: инструментируй то, что тихо деградирует, до того, как оно тебя разбудит.

✅ Exactly-Once & Fault Tolerance✅ Exactly-Once и отказоустойчивость

How exactly-once worksКак работает exactly-once

  • Exactly-once = state consistency, NOT "each message on the wire once." Every record affects state exactly once, even across failures.
  • Mechanism: asynchronous barrier snapshotting (Chandy-Lamport variant).
  • Coordinator injects checkpoint barriers into sources; barriers flow with data; on barrier, operator snapshots state.
  • Barrier alignment: multi-input op waits for the barrier on ALL inputs (buffering fast ones) before snapshotting → this is what guarantees exactly-once. At-least-once skips alignment → faster but may double-count.
  • Unaligned checkpoints: let barriers overtake in-flight data under backpressure (stores in-flight buffers in snapshot) → fast checkpoints when backpressured.
  • Exactly-once = согласованность состояния, НЕ «каждое сообщение на проводе один раз». Каждая запись влияет на состояние ровно один раз, даже при сбоях.
  • Механизм: асинхронное барьерное снапшотирование (вариант Chandy-Lamport).
  • Координатор инжектит checkpoint barriers в источники; барьеры текут с данными; при барьере оператор делает снапшот состояния.
  • Barrier alignment: оператор с несколькими входами ждёт барьера на ВСЕХ входах (буферизуя быстрые) перед снапшотом → это и гарантирует exactly-once. At-least-once пропускает alignment → быстрее, но может задвоить.
  • Unaligned checkpoints: пропускают барьеры обгонять in-flight данные при backpressure (сохраняют in-flight буферы в снапшоте) → быстрые checkpoints при backpressure.
END-TO-END EXACTLY-ONCE = consistent snapshot + cooperating source + sink [Kafka source]──barrier──▶[stateful op]──barrier──▶[2PC / idempotent sink] resettable snapshot state pre-commit on checkpoint, offsets + offsets commit on checkpoint-complete
СКВОЗНОЙ EXACTLY-ONCE = согласованный снапшот + кооперирующий source + sink [Kafka source]──barrier──▶[stateful op]──barrier──▶[2PC / идемп. sink] сбрасываемые снапшот состояния pre-commit при checkpoint, оффсеты + оффсеты commit при завершении checkpoint

Guarantee ladder & costЛестница гарантий и цена

LevelHowCost
at-most-onceno checkpointsmay lose data
at-least-oncecheckpoints, no alignmentpossible dupes
exactly-once (state)aligned checkpointsalignment latency
end-to-end EO+ transactional/idempotent sink2PC complexity
УровеньКакЦена
at-most-onceнет checkpointsможет потерять данные
at-least-oncecheckpoints, без alignmentвозможны дубли
exactly-once (state)aligned checkpointsзадержка alignment
сквозной EO+ транзакционный/идемпотентный sinkсложность 2PC
PROBEВОПРОС

"What breaks exactly-once?" → a non-transactional, non-idempotent sink. You can checkpoint perfectly and still double-write to the external system on replay.

«Что ломает exactly-once?» → не транзакционный, не идемпотентный sink. Можно идеально чекпоинтить и всё равно задвоить запись во внешнюю систему при повторе.

ANCHORЯКОРЬ

Same reasoning as Kafka offset management at IU Group: only advance the committed offset once the downstream write is durable, else gaps or dupes on restart. Flink formalizes that across the whole operator graph.

То же рассуждение, что при управлении оффсетами Kafka в IU Group: двигать закоммиченный оффсет только когда downstream-запись надёжна, иначе пропуски или дубли при рестарте. Flink формализует это на весь граф операторов.

🏗 Architecture, Backpressure & Operators🏗 Архитектура, Backpressure и операторы

Runtime componentsRuntime-компоненты

  • JobManager (master): schedules, coordinates checkpoints, handles failover. Contains Dispatcher, ResourceManager, per-job JobMaster.
  • TaskManagers (workers): run subtasks in task slots (resource isolation). Slot sharing lets a pipeline slice share a slot.
  • Dataflow: program → logical JobGraph → physical ExecutionGraph (parallel subtasks). Operators fuse into operator chains (one thread, no serialization → big perf win).
  • JobManager (мастер): планирует, координирует checkpoints, обрабатывает failover. Содержит Dispatcher, ResourceManager, per-job JobMaster.
  • TaskManagers (воркеры): выполняют subtasks в task slots (изоляция ресурсов). Slot sharing позволяет срезу пайплайна делить slot.
  • Dataflow: программа → логический JobGraph → физический ExecutionGraph (параллельные subtasks). Операторы сливаются в operator chains (один поток, без сериализации → большой выигрыш в производительности).

BackpressureBackpressure

  • Downstream can't keep up → it slows upstream via credit-based flow control (receiver grants buffer credits) instead of dropping/OOM. Pressure propagates back to source.
  • Diagnose in Web UI: bottleneck = first operator that is busy but not backpressured.
  • Causes: key skew, slow sink, GC, undersized parallelism. Heavy backpressure delays barriers → slow checkpoints (→ unaligned helps).
  • Downstream не успевает → замедляет upstream через credit-based flow control (получатель выдаёт кредиты на буфер) вместо сброса/OOM. Давление распространяется к источнику.
  • Диагностика в Web UI: узкое место = первый оператор, который занят, но не под backpressure.
  • Причины: перекос по ключам, медленный sink, GC, недостаточный parallelism. Сильный backpressure задерживает барьеры → медленные checkpoints (→ unaligned помогает).

Partitioning & low-level APIsПартиционирование и низкоуровневые API

  • keyBy hash partition (same key→same subtask) · rebalance round-robin (fix skew) · rescale local round-robin · broadcast to all (broadcast-state pattern for dynamic rules) · forward 1:1 (enables chaining).
  • KeyedProcessFunction: lowest-level — events + keyed state + event/processing-time timers. For custom state machines, dedup, timeouts, joins.
  • FlinkCEP: pattern detection ("A then B within 10m, not preceded by C") via NFA — fraud/abuse/alerting.
  • Stream joins: interval join (bounded), CoProcessFunction (manual TTL), temporal table join (as-of dimension). Always bound join state.
  • keyBy хеш-партиционирование (тот же ключ→тот же subtask) · rebalance round-robin (лечить перекос) · rescale локальный round-robin · broadcast всем (broadcast-state паттерн для динамических правил) · forward 1:1 (позволяет chaining).
  • KeyedProcessFunction: самый низкоуровневый — события + keyed state + event/processing-time таймеры. Для кастомных state machines, dedup, таймаутов, join.
  • FlinkCEP: обнаружение паттернов («A затем B в пределах 10m, не предшествует C») через NFA — фрод/абьюз/алертинг.
  • Stream joins: interval join (ограниченный), CoProcessFunction (ручной TTL), temporal table join (as-of измерение). Всегда ограничивай join-состояние.
ANCHORЯКОРЬ

Abuse/integrity detection at Meta is full of "sequence of actions within a window" logic — CEP and KeyedProcessFunction are the streaming-native way to express the patterns I built in batch.

Обнаружение абьюза/целостности в Meta полно логики «последовательность действий в окне» — CEP и KeyedProcessFunction — это streaming-native способ выразить паттерны, которые я строил в батче.

⚔ Flink vs Spark Structured Streaming · Iceberg⚔ Flink vs Spark Structured Streaming · Iceberg

DimensionFlinkSpark Structured Streaming
Modeltrue record-at-a-timemicro-batch (Continuous is limited)
Latencysingle-digit ms~100s ms to seconds
Large statefirst-class, RocksDBsupported, historically less mature
Event timemature, fine-grainedsimpler model
Batch + MLweaker batch storydominant ecosystem, MLlib
Best fitlow-latency, complex stateful, CEP, alertsbatch+stream on existing Spark, seconds OK
ИзмерениеFlinkSpark Structured Streaming
Модельнастоящий record-at-a-timemicro-batch (Continuous ограничен)
Задержкаединицы мс~100 мс - секунды
Большое состояниеfirst-class, RocksDBподдерживается, исторически менее зрелое
Event timeзрелое, тонко настраиваемоеболее простая модель
Batch + MLслабее по батчудоминирующая экосистема, MLlib
Лучше подходитнизкая задержка, сложное stateful, CEP, алертыбатч+поток на существующем Spark, секунды OK

One-liner: Flink for true low-latency, large-stateful, event-time-correct streaming. Spark SS when seconds are fine and you want one engine for batch + stream + ML on an existing stack.

Одной строкой: Flink для настоящего низкозадержечного, large-stateful, event-time-корректного стриминга. Spark SS, когда секунды приемлемы и нужен один движок для батча + потока + ML на существующем стеке.

Iceberg + the canonical pipelineIceberg и канонический пайплайн

Kafka ──▶ Flink ──────────▶ Iceberg ──▶ Trino / Spark (exactly-once, (ACID snapshots, (analytics, stateful, EO) schema evolution, time travel) hidden partitioning)
Kafka ──▶ Flink ──────────▶ Iceberg ──▶ Trino / Spark (exactly-once, (ACID снапшоты, (аналитика, stateful, EO) schema evolution, time travel) скрытое партиционирование)
  • Flink reads/writes Iceberg (stream + batch); streaming writes commit as Iceberg snapshots aligned with Flink checkpoints → exactly-once into a lakehouse.
  • Iceberg = open table format giving ACID snapshot isolation, schema evolution, time travel over object storage. Open-source analog to Snowflake/BigQuery/Hive formats you've used.
  • Flink читает/пишет Iceberg (поток + батч); потоковые записи коммитятся как Iceberg-снапшоты, выровненные с Flink checkpoints → exactly-once в lakehouse.
  • Iceberg = открытый табличный формат, дающий ACID snapshot isolation, schema evolution, time travel поверх объектного хранилища. Открытый аналог форматов Snowflake/BigQuery/Hive, которые ты использовал.
GAP — say it cleanlyПРОБЕЛ — скажи чётко

"Iceberg is a study area for me; I know it as the ACID, schema-evolving, time-travel table format that unifies realtime + historical. The pipeline pairs Flink + Kafka + Iceberg and that Kafka→Flink→Iceberg→Trino pattern is exactly the architecture I'd draw."

«Iceberg — зона изучения для меня; я знаю его как ACID, schema-evolving, time-travel табличный формат, объединяющий realtime + исторические данные. Пайплайн связывает Flink + Kafka + Iceberg, и паттерн Kafka→Flink→Iceberg→Trino — это именно та архитектура, которую я бы нарисовал.»

🃏 Flip-to-Reveal Self-Quiz🃏 Самопроверка с переворачиванием

Tap a card to flip. Answer out loud first.

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

What is a watermark in one sentence?Что такое watermark одним предложением?
tap to revealнажми, чтобы открыть
A marker asserting "no event with timestamp ≤ t will arrive after this." It tracks event-time progress and fires event-time windows. Operator watermark = min across inputs.Маркер, утверждающий «событий с timestamp ≤ t после этого не будет». Отслеживает прогресс event-time и запускает event-time окна. Watermark оператора = min по входам.
How does Flink achieve exactly-once?Как Flink достигает exactly-once?
tap to revealнажми, чтобы открыть
Aligned checkpoint barriers (Chandy-Lamport snapshots): operators wait for the barrier on all inputs, snapshot state consistently. End-to-end needs a transactional/idempotent sink + resettable source.Выровненные checkpoint barriers (Chandy-Lamport снапшоты): операторы ждут барьера на всех входах, делают снапшот состояния согласованно. Сквозной требует транзакционного/идемпотентного sink + сбрасываемого source.
Checkpoint vs Savepoint?Checkpoint vs Savepoint?
tap to revealнажми, чтобы открыть
Checkpoint = automatic, Flink-managed, for crash recovery. Savepoint = manual, user-owned, portable — for code upgrades, rescale, version/cluster migration.Checkpoint = автоматический, управляется Flink, для восстановления после падения. Savepoint = ручной, в собственности пользователя, портируемый — для апгрейдов кода, масштабирования, миграции версии/кластера.
Why set .uid() on operators?Зачем ставить .uid() на операторах?
tap to revealнажми, чтобы открыть
On restore Flink maps saved state to operators by ID. Auto IDs change when topology changes → restore fails / state dropped. Stable UIDs let state survive code refactors and rescale.При восстановлении Flink сопоставляет сохранённое состояние с операторами по ID. Авто-ID меняются при изменении топологии → восстановление падает / состояние теряется. Стабильные UID позволяют состоянию пережить рефакторинг кода и масштабирование.
HashMap vs RocksDB state backend?HashMap vs RocksDB state backend?
tap to revealнажми, чтобы открыть
HashMap = on JVM heap, fast, small state. EmbeddedRocksDB = on local disk, serialized, supports state bigger than memory + incremental checkpoints, slightly higher access latency.HashMap = на JVM heap, быстрое, маленькое состояние. EmbeddedRocksDB = на локальном диске, сериализованное, поддерживает состояние больше памяти + инкрементальные checkpoints, чуть выше задержка доступа.
Tumbling vs sliding vs session window?Tumbling vs sliding vs session окно?
tap to revealнажми, чтобы открыть
Tumbling = fixed, non-overlapping. Sliding = fixed size + slide, overlapping. Session = dynamic, closes after a gap of inactivity (sessionization).Tumbling = фиксированное, без перекрытия. Sliding = фиксированный размер + слайд, с перекрытием. Session = динамическое, закрывается после паузы в активности (сессионизация).
What is barrier alignment and its cost?Что такое barrier alignment и его цена?
tap to revealнажми, чтобы открыть
A multi-input operator waits for the checkpoint barrier on ALL inputs (buffering faster ones) before snapshotting. Guarantees exactly-once; costs latency. Skipping it = at-least-once. Unaligned checkpoints relax it under backpressure.Оператор с несколькими входами ждёт checkpoint барьера на ВСЕХ входах (буферизуя быстрые) перед снапшотом. Гарантирует exactly-once; стоит задержку. Пропуск = at-least-once. Unaligned checkpoints ослабляют при backpressure.
How does Flink handle backpressure?Как Flink обрабатывает backpressure?
tap to revealнажми, чтобы открыть
Credit-based flow control: receivers grant buffer credits; senders transmit only with credits. Pressure propagates to the source — no drops/OOM. Diagnose in Web UI (busy vs backpressured).Credit-based flow control: получатели выдают кредиты на буфер; отправители передают только при наличии кредитов. Давление распространяется к источнику — нет сброса/OOM. Диагностика в Web UI (занят vs под backpressure).
How do you avoid unbounded state?Как избежать неограниченного состояния?
tap to revealнажми, чтобы открыть
State TTL (StateTtlConfig) to expire old keys; bound stream joins via interval/temporal; RocksDB incremental checkpoints. Otherwise recovery slows and disk blows up.State TTL (StateTtlConfig) для устаревания старых ключей; ограничь stream join через interval/temporal; RocksDB инкрементальные checkpoints. Иначе восстановление замедляется и диск взрывается.
Key groups — what and why care?Key groups — что и зачем заботиться?
tap to revealнажми, чтобы открыть
Atomic unit of keyed-state redistribution on rescale. Max parallelism fixes their count, capping rescale without state migration — so set it deliberately at job creation.Атомарная единица перераспределения keyed-state при масштабировании. Max parallelism фиксирует их число, ограничивая масштабирование без миграции состояния — поэтому устанавливай осознанно при создании джоба.

⚠ Gotchas & Failure-Mode Drills⚠ Подводные камни и разбор режимов сбоев

Top production gotchasГлавные подводные камни в продакшене

  • No .uid() → can't restore from savepoint after a code edit.
  • No state TTL → state grows forever, recovery gets slower over time.
  • Idle Kafka partition → min watermark stalls → windows never fire. Use withIdleness.
  • Non-idempotent sink → exactly-once is a lie; you double-write on replay.
  • Key skew → one subtask backpressured → checkpoints time out. Salt / two-phase aggregate.
  • Dropping late data silently → use side outputs to capture it.
  • Unbounded stream join → state explosion. Use interval/temporal joins or TTL.
  • Нет .uid() → не восстановить из savepoint после правки кода.
  • Нет state TTL → состояние растёт вечно, восстановление замедляется со временем.
  • Idle партиция Kafka → min watermark застревает → окна не запускаются. Использовать withIdleness.
  • Не идемпотентный sink → exactly-once — ложь; ты задваиваешь запись при повторе.
  • Перекос по ключам → один subtask под backpressure → checkpoints таймаутятся. Salt / двухфазная агрегация.
  • Тихий сброс опоздавших данных → использовать side outputs, чтобы захватить их.
  • Неограниченный stream join → взрыв состояния. Использовать interval/temporal join или TTL.

Drill: "Checkpoints keep timing out under load."Разбор: «Checkpoints таймаутятся под нагрузкой»

  1. Web UI: is it alignment time (backpressure) or snapshot time (state size)?
  2. Alignment-bound → find bottleneck (busy, not backpressured); fix skew / scale / speed sink; try unaligned checkpoints.
  3. State-size-bound → incremental checkpoints (RocksDB) + state TTL; relax interval/timeout.
  4. Sink-bound (slow 2PC commit) → investigate external latency; idempotent sink may beat 2PC.
  5. Confirm with metrics, not vibes.
  1. Web UI: это alignment time (backpressure) или snapshot time (размер состояния)?
  2. Ограничен alignment → найти узкое место (занят, но не под backpressure); чинить перекос / масштабировать / ускорить sink; попробовать unaligned checkpoints.
  3. Ограничен размером состояния → инкрементальные checkpoints (RocksDB) + state TTL; ослабить интервал/таймаут.
  4. Ограничен sink (медленный 2PC commit) → исследовать внешнюю задержку; идемпотентный sink может быть лучше 2PC.
  5. Подтверждать метриками, а не ощущениями.

When NOT to use FlinkКогда НЕ использовать Flink

  • Pure batch ETL where seconds-minutes latency is fine and the org is on Spark.
  • Low volume where a consumer or cron suffices — Flink's ops overhead isn't justified.
  • No JVM/streaming ops capacity — exactly-once + state + savepoint discipline is real cost.
  • Picking boring batch when realtime isn't needed is senior judgment.
  • Чистый батч ETL, где задержка в секундах-минутах приемлема, и организация на Spark.
  • Малый объём, где хватает консьюмера или cron — операционные расходы Flink не оправданы.
  • Нет JVM/streaming-ops компетенций — exactly-once + состояние + savepoint дисциплина — реальная цена.
  • Выбирать скучный батч, когда realtime не нужен — это senior-суждение.
CLOSING ANCHOR (gap → strength)ЗАКЛЮЧИТЕЛЬНЫЙ ЯКОРЬ (пробел → сила)

"No prod Flink yet, but I've shipped the hard parts under other names: Kafka exactly-once offset reasoning, Scala/Spark stateful pipelines, Airflow migrations saving €1M+. Day one I'd enforce UIDs, state TTL, deterministic event-time logic, MiniCluster integration tests and checkpoint monitoring — the craftsmanship and failure-mode rigor needed."

«Flink в проде пока нет, но я запустил сложные части под другими именами: Kafka exactly-once рассуждения про оффсеты, Scala/Spark stateful-пайплайны, Airflow-миграции, сэкономившие €1M+. С первого дня я бы внедрил UID, state TTL, детерминированную event-time логику, MiniCluster integration tests и мониторинг checkpoints — необходимые мастерство и строгость в режимах сбоев.»