keyBy; rescaled via key groups; bound it with TTL.keyBy; масштабируется через key groups; ограничь через TTL."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 — это те же рассуждения о режимах сбоев, которые я уже делал на тех системах, поэтому дай мне порассуждать с первых принципов.»
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.
| Type | Meaning | Trade-off |
|---|---|---|
| Event | timestamp in the record | correct + replayable; needs watermarks → latency |
| Ingestion | assigned at source op | middle ground |
| Processing | operator wall-clock | fastest; non-deterministic, wrong under skew |
| Тип | Смысл | Компромисс |
|---|---|---|
| Event | метка времени в записи | корректно + повторяемо; требует watermarks → задержка |
| Ingestion | назначается в source-операторе | золотая середина |
| Processing | wall-clock оператора | быстрее всего; не детерминировано, неверно при перекосе |
"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-ы/повторы дают идентичные результаты, потому что логика использует встроенную метку времени, а не момент запуска джоба.
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.
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
withIdleness() marks it idle.withIdleness() помечает его idle."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.
| Window | Shape | Use for |
|---|---|---|
| Tumbling | fixed size, no overlap, 1 event → 1 window | periodic buckets ("count per 1 min") |
| Sliding | size + slide; overlap; event in many windows | moving averages ("5-min count every 1 min") |
| Session | dynamic, closes after gap of inactivity | user sessionization ("clicks until 30-min idle") |
| Global | one window per key, custom Trigger fires it | building 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());
OutputTag): route too-late events to a separate stream instead of losing them.OutputTag): направить слишком поздние события в отдельный поток вместо потери.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.
keyBy): one independent instance per key. Types: ValueState, ListState, MapState, ReducingState, AggregatingState.keyBy): один независимый экземпляр на ключ. Типы: ValueState, ListState, MapState, ReducingState, AggregatingState.| Backend | Where | When |
|---|---|---|
| HashMap | JVM heap | small/fast state, low latency |
| EmbeddedRocksDB | local disk, off-heap, serialized | huge state > memory; incremental checkpoints |
| Backend | Где | Когда |
|---|---|---|
| HashMap | JVM heap | маленькое/быстрое состояние, низкая задержка |
| EmbeddedRocksDB | локальный диск, off-heap, сериализованное | огромное состояние > памяти; инкрементальные checkpoints |
Checkpoint storage (HDFS/S3) is separate: it's where the snapshot lands.
Checkpoint storage (HDFS/S3) — отдельно: это куда попадает снапшот.
StateTtlConfig): expire old keys or state grows forever → slow recovery, disk blowup. Essential at per-user, billions-of-keys scale.StateTtlConfig): устаревание старых ключей, иначе состояние растёт вечно → медленное восстановление, взрыв диска. Критично на масштабе миллиардов ключей по пользователям."What if state grows unbounded?" → State TTL + RocksDB incremental checkpoints; bound joins with intervals.
«Что, если состояние растёт неограниченно?» → State TTL + RocksDB инкрементальные checkpoints; ограничь join-ы интервалами.
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 тренировал тот же инстинкт: ключить данные, держать состояние на ключ ограниченным, детерминированное восстановление.
| Checkpoint | Savepoint | |
|---|---|---|
| Purpose | automatic fault tolerance | manual operational snapshot |
| Trigger | periodic, by Flink | on-demand, by you/CI |
| Ownership | Flink-managed, may be cleaned up | user-owned, retained, portable |
| Use case | recover from crash | upgrade code, rescale, migrate cluster/version |
| Checkpoint | Savepoint | |
|---|---|---|
| Цель | автоматическая отказоустойчивость | ручной операционный снапшот |
| Триггер | периодически, 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 совместимы.
"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-операцию с первого дня.
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: инструментируй то, что тихо деградирует, до того, как оно тебя разбудит.
| Level | How | Cost |
|---|---|---|
| at-most-once | no checkpoints | may lose data |
| at-least-once | checkpoints, no alignment | possible dupes |
| exactly-once (state) | aligned checkpoints | alignment latency |
| end-to-end EO | + transactional/idempotent sink | 2PC complexity |
| Уровень | Как | Цена |
|---|---|---|
| at-most-once | нет checkpoints | может потерять данные |
| at-least-once | checkpoints, без alignment | возможны дубли |
| exactly-once (state) | aligned checkpoints | задержка alignment |
| сквозной EO | + транзакционный/идемпотентный sink | сложность 2PC |
"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. Можно идеально чекпоинтить и всё равно задвоить запись во внешнюю систему при повторе.
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 формализует это на весь граф операторов.
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).keyBy хеш-партиционирование (тот же ключ→тот же subtask) · rebalance round-robin (лечить перекос) · rescale локальный round-robin · broadcast всем (broadcast-state паттерн для динамических правил) · forward 1:1 (позволяет chaining).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 способ выразить паттерны, которые я строил в батче.
| Dimension | Flink | Spark Structured Streaming |
|---|---|---|
| Model | true record-at-a-time | micro-batch (Continuous is limited) |
| Latency | single-digit ms | ~100s ms to seconds |
| Large state | first-class, RocksDB | supported, historically less mature |
| Event time | mature, fine-grained | simpler model |
| Batch + ML | weaker batch story | dominant ecosystem, MLlib |
| Best fit | low-latency, complex stateful, CEP, alerts | batch+stream on existing Spark, seconds OK |
| Измерение | Flink | Spark Structured Streaming |
|---|---|---|
| Модель | настоящий record-at-a-time | micro-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 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 — это именно та архитектура, которую я бы нарисовал.»
Tap a card to flip. Answer out loud first.
Нажми на карточку, чтобы перевернуть. Сначала ответь вслух.
.uid() on operators?Зачем ставить .uid() на операторах?.uid() → can't restore from savepoint after a code edit.withIdleness..uid() → не восстановить из savepoint после правки кода.withIdleness."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 — необходимые мастерство и строгость в режимах сбоев.»