Data Pipeline Design & Workflow Orchestration — DE Interview PrepДизайн пайплайнов и оркестрация — Подготовка к DE-интервью

Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE. Шпаргалка для быстрого повторения перед Data Engineering интервью. Заточена под Senior/Lead DE.
batch vs streamingбатч vs стриминг idempotency & backfillsидемпотентность и backfills Airflow deep-diveAirflow детально time-sensitive releasesрелизы с жёсткими SLA revisions & point-in-timeревизии и point-in-time dbt incremental + snapshotsdbt инкрементальные + снапшоты

Click any section header to collapse. 🌐 RU / 🌙 Dark toggles (saved). Flip cards at the bottom for self-quiz. Hooks tie answers to McMakler / Meta / IU / Front Tier. Кликай на заголовок секции, чтобы свернуть. 🌐 RU / 🌙 Dark переключают (сохраняется). Карточки внизу — самопроверка. Якоря привязывают ответы к McMakler / Meta / IU / Front Tier.

1 · Fundamentals & Vocabulary1 · Основы и терминология

Core definitionsКлючевые определения

  • Idempotency — re-running a task with the same input yields the same end state; no dupes, no drift.
  • Backfill — re-run historical partitions to (re)populate or correct past data.
  • Reprocessing — recompute affected data after a bug fix or a data revision.
  • Late-arriving data — a record for an old event_time that lands now.
  • Revision / vintage — a new value for an already-published observation (CPI/GDP get revised for months).
  • Point-in-time / as-of — what was known at a past moment, before later revisions (no look-ahead bias).
  • Watermark — boundary marking how complete a stream is up to a given event time.
  • Идемпотентность (Idempotency) — повторный запуск задачи с теми же входными данными даёт то же конечное состояние; никаких дублей, никакого дрифта.
  • Backfill — повторный запуск исторических партиций для (пере)наполнения или исправления прошлых данных.
  • Reprocessing — перевычисление затронутых данных после багфикса или ревизии данных.
  • Late-arriving data — запись для старого event_time, которая прилетает сейчас.
  • Revision / vintage — новое значение для уже опубликованного наблюдения (CPI/GDP пересматриваются месяцами).
  • Point-in-time / as-of — что было известно в прошлый момент, до последующих ревизий (никакого взгляда в будущее).
  • Watermark — граница, отмечающая, насколько полон поток до заданного event time.

Batch vs StreamingБатч vs стриминг

BatchStreaming
Databounded, periodicunbounded, continuous
Latencyminutes–hourssub-second–seconds
Replayeasy (re-run partition)harder (offsets, watermarks)
Fitmost macro releasesticks, clickstream, LLM events
БатчСтриминг
Данныеограниченные, периодичныебесконечные, непрерывные
Задержкаминуты–часыдоли секунды–секунды
Повторлегко (перезапуск партиции)сложнее (оффсеты, watermarks)
Подходит длябольшинства макрорелизовтики, кликстрим, события LLM
hookBatch backbone at McMakler (Airflow); streaming at IU/Syntea (Kafka, 80k+ students). Choose per-feed by latency & arrival pattern, not religion.
якорьОснова на батче в McMakler (Airflow); стриминг в IU/Syntea (Kafka, 80k+ студентов). Выбирай по фиду на основе задержки и паттерна прибытия данных, а не по религии.

The layered pattern everything reduces toСлоёный паттерн, к которому всё сводится

RAW (bronze) STAGING (silver) MARTS (gold) SERVE immutable land -> typed + normalized -> conformed model -> query layer / BI partitioned by quarantine bad rows long format, point-in-time / as-of views load_date one row/record revisions kept
RAW (bronze) STAGING (silver) MARTS (gold) SERVE неизменная посадка typed + normalized -> conформленная модель query layer / BI партиции по карантин плохих строк long format, point-in-time / as-of виды load_date одна строка/запись ревизии сохранены

Raw is never mutated — it is your replay source. Vendor quirks die in staging. One canonical model in marts.

Raw никогда не мутируется — это источник для повтора. Причуды вендоров умирают в staging. Одна каноничная модель в marts.

2 · Whiteboard: End-to-End Vendor Macro Feed → Validated Dataset2 · Whiteboard: E2E фид от вендора → валидированный датасет

interviewer will probeThis is the #1 design prompt. Prepare it cold. They will push on the validation gate and on revisions.
спросит интервьюерЭто промпт дизайна №1. Готовь с холодного старта. Будут давить на validation gate и на revisions.

Say it as a DAG in this orderРассказывай как DAG в таком порядке

[release calendar] (CPI prints 13:30 UK on known date) | v (1) TRIGGER ----> event/push (SFTP, webhook, Kafka) + deferrable sensor safety-net | (2) LAND ------> object store, immutable, partition vendor/dataset/load_date + hash | (3) PARSE -----> typed staging; normalize dates/decimals/units; quarantine bad rows | (4) VALIDATE GATE ==X==> bad data NEVER published (schema, counts, freshness, | nulls/range, PK uniqueness, anomaly) [fail closed] v (5) MODEL -----> canonical: (series_id, observation_date, value, revision_id, load_ts) | derived (MoM/YoY/rolling) = downstream models, not in-place v (6) SERVE -----> warehouse + query layer; as-of views cross-cut: orchestration (Airflow) · idempotency (partition overwrite) · observability (freshness/completeness SLA) · lineage · alert+runbook
[календарь релизов] (CPI выходит 13:30 UK в известную дату) | v (1) TRIGGER ----> event/push (SFTP, webhook, Kafka) + deferrable sensor как сеть безопасности | (2) LAND ------> object store, неизменный, партиция vendor/dataset/load_date + хеш | (3) PARSE -----> типизированный staging; нормализация дат/decimal/units; карантин плохих строк | (4) VALIDATE GATE ==X==> плохие данные НИКОГДА не публикуются (схема, counts, freshness, | null/диапазон, PK uniqueness, аномалия) [fail closed] v (5) MODEL -----> каноничный: (series_id, observation_date, value, revision_id, load_ts) | производные (MoM/YoY/rolling) = downstream модели, не in-place v (6) SERVE -----> warehouse + query layer; as-of виды cross-cut: оркестрация (Airflow) · идемпотентность (partition overwrite) · наблюдаемость (freshness/completeness SLA) · lineage · alert+runbook
hook"Same shape as the Meta 10-logging-system consolidation: profile heterogeneous sources, land immutably, conform to one model, treat downstream as first-class with validation before publish. There I cut 47% redundant volume at ~3.6B events/day."
якорь«Та же форма, что и Meta консолидация 10 систем логирования: профилировать гетерогенные источники, посадить неизменно, привести к одной модели, обращаться с downstream как с first-class с валидацией перед публикацией. Там я срезал 47% избыточного объёма на ~3,6B событий/день.»

Validation gate — the heart of the answerValidation gate — сердце ответа

  • Schema/contract: columns, types, enums; detect drift
  • Completeness: row count vs expected; all series/dates present (gap detection)
  • Freshness: landed within SLA of the release time
  • Validity: non-null required fields; value in plausible range; unit sanity
  • Uniqueness: (series_id, observation_date, revision) unique
  • Anomaly: N-sigma / sudden jump (catches a misplaced decimal)
  • Схема/контракт: колонки, типы, enums; детект дрифта
  • Полнота: число строк vs ожидаемое; все серии/даты присутствуют (детект пропусков)
  • Свежесть: приземлилось в пределах SLA от времени релиза
  • Валидность: обязательные поля не null; значение в правдоподобном диапазоне; вменяемость unit
  • Уникальность: (series_id, observation_date, revision) уникальны
  • Аномалия: N-sigma / резкий скачок (ловит decimal, сдвинутую запятую)

Behavior: quarantine bad rows; fail closed on systemic problems; alert with runbook; remediation is first-class (re-pull, re-run partition).

Поведение: плохие строки в карантин; fail closed при системных проблемах; алерт с runbook; восстановление — first-class (перезапрос, перезапуск партиции).

Serving correctlyПравильная подача

  • Keep all vintages; "current" value = latest revision per (series_id, observation_date) as a view.
  • As-of view = revision known on or before date X -> no look-ahead bias.
  • Query layer reads the conformed marts, not staging.
  • Хранить все vintage; «текущее» значение = последняя ревизия на (series_id, observation_date) как view.
  • As-of view = ревизия, известная на дату X или раньше -> нет взгляда в будущее.
  • Query layer читает конформленные marts, а не staging.
gotchaNever overwrite a published value. A revision is an append of a new vintage, not an UPDATE.
gotchaНикогда не перезаписывай опубликованное значение. Ревизия — это append нового vintage, а не UPDATE.

3 · Idempotency · Backfills · Late Data · Reprocessing3 · Идемпотентность · Backfills · Поздние данные · Reprocessing

Idempotency: howИдемпотентность: как

  • Partition overwrite / delete+insert keyed by deterministic partition (load_date / observation_date): re-run replaces, never appends.
  • MERGE / upsert on a stable natural key, not blind INSERT.
  • Pass the logical date in ({{ ds }} / logical_date); no now()/random in business logic.
  • Streaming: exactly-once sink or dedup on content hash.
  • Partition overwrite / delete+insert по детерминистичной партиции (load_date / observation_date): повтор заменяет, а не добавляет.
  • MERGE / upsert по стабильному натуральному ключу, а не слепой INSERT.
  • Передавать logical date ({{ ds }} / logical_date); никаких now()/random в бизнес-логике.
  • Стриминг: exactly-once sink или dedup по content hash.
Why backfills need it: a backfill re-runs N partitions in any order, possibly in parallel. Non-idempotent tasks double-count. Idempotency = safe re-run of any date range.
Зачем backfills это нужно: backfill перезапускает N партиций в любом порядке, возможно параллельно. Не-идемпотентные задачи считают дважды. Идемпотентность = безопасный повтор любого диапазона дат.

Late & revised dataПоздние и пересмотренные данные

  • Separate event_time from processing_time.
  • Partition by observation_date so a late record lands in the correct historical partition.
  • A late value = a new revision keyed to that date; re-run that one partition idempotently.
  • Recompute derived series (MoM/YoY/rolling) for affected windows — one revision shifts several points.
  • Streaming: watermark + grace window, then finalize.
  • Разделять event_time и processing_time.
  • Партиция по observation_date, чтобы поздняя запись попала в правильную историческую партицию.
  • Позднее значение = новая ревизия с ключом на эту дату; перезапуск этой одной партиции идемпотентно.
  • Перевычислить производные ряды (MoM/YoY/rolling) для затронутых окон — одна ревизия сдвигает несколько точек.
  • Стриминг: watermark + grace window, затем финализация.
hook"Fixing the Meta entity-resolution defect across 6.2B records, idempotent partition-level reprocessing is what made the correction safe at scale and restored 100+ downstream metrics."
якорь«Исправляя Meta баг entity-resolution на 6,2B записей, идемпотентный reprocessing на уровне партиций — это то, что сделало исправление безопасным на масштабе и восстановило 100+ downstream метрик.»
WRONG (append) RIGHT (idempotent overwrite) run1: INSERT day=10 -> rows run1: OVERWRITE day=10 -> state S run2: INSERT day=10 -> DUP rows X run2: OVERWRITE day=10 -> state S (same) backfill day=1..30 in any order -> converges
НЕПРАВИЛЬНО (append) ПРАВИЛЬНО (idempotent overwrite) run1: INSERT day=10 -> строки run1: OVERWRITE day=10 -> состояние S run2: INSERT day=10 -> ДУБЛИ X run2: OVERWRITE day=10 -> состояние S (то же) backfill day=1..30 в любом порядке -> сходится

4 · Airflow Deep-Dive (candidate strength)4 · Airflow детально (сильная сторона кандидата)

DAGs & dependenciesDAG-и и зависимости

  • DAG = directed acyclic graph of tasks; edges (>>) set ordering; scheduler runs a task when upstreams satisfy its trigger_rule.
  • Keep tasks small, single-purpose, idempotent.
  • Group with TaskGroups; share logic via custom operators/hooks; avoid heavy top-level code (runs on every parse).
  • Cross-DAG: prefer Datasets (data-aware scheduling) or ExternalTaskSensor over timing assumptions.
  • DAG = направленный ациклический граф задач; рёбра (>>) задают порядок; шедулер запускает задачу, когда upstream удовлетворяют её trigger_rule.
  • Держать задачи маленькими, single-purpose, идемпотентными.
  • Группировать через TaskGroups; делить логику через custom операторы/хуки; избегать тяжёлого top-level кода (выполняется на каждом parse).
  • Cross-DAG: предпочитать Datasets (data-aware шедулинг) или ExternalTaskSensor поверх timing допущений.

SensorsСенсоры

  • Wait for a condition (file landed, partition ready, external task done).
  • poke holds a worker slot the whole wait -> can starve the pool.
  • reschedule releases the slot between checks -> use for any wait of minutes+.
  • Deferrable operators / triggerer = async wait, zero worker slot. Preferred for "wait for the CPI file".
  • Always set a timeout.
  • Ожидать условие (файл приземлился, партиция готова, внешняя задача завершена).
  • poke держит worker slot всё время ожидания -> может заморить pool голодом.
  • reschedule освобождает slot между проверками -> использовать для ожидания минут+.
  • Deferrable операторы / triggerer = async ожидание, нулевой worker slot. Предпочтительно для «ждать файл CPI».
  • Всегда ставить timeout.
gotchaA fleet of poke sensors is the classic way to deadlock a Celery pool.
gotchaФлот poke сенсоров — классический способ заблокировать Celery pool.

Retries · SLAs · AlertingRetries · SLA · Алертинг

  • Retries: retries, retry_delay, retry_exponential_backoff. Retry transient (network/API blip); never mask deterministic logic bugs.
  • SLA: fires sla_miss_callback if not done within duration after logical time ("land by 13:45"). Known quirks -> many teams build custom freshness checks.
  • Alerting: on_failure_callback -> Slack/PagerDuty + runbook; dedupe + severity-tier to beat alert fatigue.
  • Retries: retries, retry_delay, retry_exponential_backoff. Повторять transient (сетевой/API сбой); никогда не маскировать детерминистичные логические баги.
  • SLA: запускает sla_miss_callback, если не завершилось в течение duration после logical time («посадка к 13:45»). Известные quirks -> многие команды строят custom freshness проверки.
  • Алертинг: on_failure_callback -> Slack/PagerDuty + runbook; dedup + severity-tier, чтобы победить alert fatigue.

Dynamic DAGs & task mappingДинамические DAG-и и task mapping

  • Dynamic DAGs: generate DAGs/tasks from a config registry (one per vendor feed). Watch parse cost.
  • Dynamic task mapping (.expand(), 2.3+): fan out over a runtime list (one task per file/partition).
  • Right for many homogeneous units; wrong when count explodes or each needs bespoke logic.
  • Динамические DAG-и: генерировать DAG-и/задачи из config registry (один на вендорский фид). Следи за стоимостью parse.
  • Dynamic task mapping (.expand(), 2.3+): развернуть по runtime списку (одна задача на файл/партицию).
  • Подходит для многих гомогенных единиц; не подходит, когда счёт взрывается или каждой нужна кастомная логика.
hook"Onboard many vendor feeds via a config-registry-generated DAG + dynamic mapping over the day's files: adding a feed is a config change, not new code."
якорь«Онбордить много вендорских фидов через DAG, сгенерированный из config-registry + динамический mapping по файлам дня: добавление фида — это изменение конфига, а не новый код.»

Executors & scaling workersExecutors и масштабирование воркеров

ExecutorModelUse whenWatch out
Sequentialone task (SQLite)local debugnever prod
Localsubprocesses, 1 hostsmall single-nodehost is ceiling
Celerydistributed workers + brokerclassic horizontal scalebroker + fleet ops
Kubernetes1 pod / taskbursty, isolationpod startup latency
CeleryKubernetesbothmixed workloadscomplexity
ExecutorМодельИспользовать когдаСледить за
Sequentialодна задача (SQLite)локальный debugникогда в prod
Localsubprocesses, 1 хостмаленький single-nodeхост — потолок
Celeryраспределённые воркеры + брокерклассический horizontal масштабброкер + fleet ops
Kubernetes1 pod / задачаbursty, изоляциязадержка старта pod
CeleryKubernetesобасмешанные нагрузкисложность

Scaling levers: parallelism, max_active_tasks_per_dag, max_active_runs_per_dag, worker count / worker_concurrency, pools (protect a vendor rate limit), priority_weight. Keep tasks light; push heavy compute to warehouse/Spark, not the worker.

Рычаги масштабирования: parallelism, max_active_tasks_per_dag, max_active_runs_per_dag, число воркеров / worker_concurrency, pools (защищать rate limit вендора), priority_weight. Держать задачи лёгкими; тяжёлые вычисления пушить в warehouse/Spark, а не в воркер.

# config-registry-driven ingestion + dynamic mapping
with DAG("vendor_feed", schedule="@daily", catchup=False) as dag:
    files = list_drop(vendor="econ_cpi")            # runtime list
    parsed = parse_one.expand(path=files)         # dynamic task mapping
    gate = DataQualityCheck(pool="vendor_api")    # pool protects rate limit
    parsed >> gate >> publish()                    # fail closed before publish
hook"At McMakler I migrated a proprietary orchestrator to Airflow: killed hundreds of K EUR/yr licensing, DAGs in version control, big hiring pool. Honest trade-off: you then operate the scheduler + metadata DB + workers, and Airflow is poor for sub-minute streaming."
якорь«В McMakler я мигрировал проприетарный оркестратор на Airflow: убил сотни тысяч EUR/год лицензий, DAG-и в version control, большой пул найма. Честный trade-off: потом ты эксплуатируешь шедулер + metadata DB + воркеров, и Airflow слаб для sub-minute стриминга.»

5 · Orchestration Alternatives & Trade-offs5 · Альтернативы оркестрации и trade-offs

ToolCore modelStrengthsWeaknesses / fit
Airflowtask DAGs, schedule-centricmature, huge ecosystem, ubiquitous, strong batch schedulingtask- not asset-centric (historically); streaming/low-latency awkward; ops weight
Dagstersoftware-defined assets, data-awareassets+lineage first-class, strong typing, great local dev, asset checks (DQ built in)smaller ecosystem; newer ops patterns
Prefectdynamic Python flows, code-firstvery Pythonic, dynamic at runtime, light to start, hybrid executionless of a batch "platform"; structure can hide in code
Temporaldurable execution, workflows-as-codedurable state, retries/timeouts as primitives, long-running stateful + service orchestrationnot a data scheduler; you build data semantics yourself
ИнструментОсновная модельСильные стороныСлабые стороны / подходит для
Airflowtask DAG-и, schedule-centricзрелый, огромная экосистема, вездесущий, сильный batch шедулингtask-, а не asset-centric (исторически); стриминг/low-latency неудобно; операционный вес
Dagstersoftware-defined assets, data-awareassets+lineage first-class, сильная типизация, отличный local dev, asset checks (DQ встроен)меньшая экосистема; новые ops паттерны
Prefectдинамические Python flows, code-firstочень Pythonic, динамика runtime, лёгкий старт, hybrid исполнениеменее batch «платформа»; структура может спрятаться в коде
Temporaldurable execution, workflows-as-codedurable state, retries/timeouts как примитивы, долгоживущие stateful + оркестрация сервисовне дата-шедулер; data семантика строится вручную
My pick: For hundreds of scheduled, validated batch feeds -> Airflow as backbone, or Dagster if I want assets + lineage + DQ native (attractive for lineage/metadata emphasis). Temporal only for long-running, stateful, human-in-the-loop workflows (e.g. a multi-step onboarding). Kafka is transport, not the orchestrator.
Мой выбор: Для сотен запланированных, валидированных батч-фидов -> Airflow как основа, или Dagster, если хочу assets + lineage + DQ native (привлекательно при упоре на lineage/metadata). Temporal только для долгих, stateful, human-in-the-loop воркфлоу (например multi-step онбординг). Kafka — это транспорт, а не оркестратор.
interviewer will probe"When is Dagster's asset model better?" -> when lineage/DQ are central: the asset (dataset) is the primitive, so dependencies, freshness policies, and asset checks are data-aware by default; the lineage graph is the orchestration graph. Airflow 2.4+ Datasets narrow the gap.
спросит интервьюер«Когда модель assets Dagster лучше?» -> когда lineage/DQ центральны: asset (датасет) — примитив, поэтому зависимости, freshness политики и asset checks data-aware по умолчанию; граф lineage и есть граф оркестрации. Airflow 2.4+ Datasets сужают разрыв.

6 · Scheduling Time-Sensitive Releases6 · Шедулинг релизов с жёсткими SLA

interviewer will probe"CPI must land minutes after the government print. How do you guarantee that?"
спросит интервьюер«CPI должен приземлиться через минуты после публикации. Как гарантировать?»
  • Release calendar as data (date + exact time, e.g. 13:30 UK) drives scheduling. Never hardcode crons.
  • Event/push trigger first; deferrable sensor as safety net from the known release time. No coarse cron.
  • Tight SLA + escalation measured from release time; miss -> page with runbook (re-pull, check vendor status).
  • Pre-warm the DAG/workers; reserve a pool so a noisy neighbor can't starve the critical feed.
  • Fast validation path: must-have checks (schema, single-row sanity, range) on the critical path; defer heavy checks.
  • Fail visible, not silent: for a client on the Terminal, late-but-correct beats fast-but-wrong.
  • Календарь релизов как данные (дата + точное время, напр. 13:30 UK) управляет шедулингом. Никогда не хардкодить cron-ы.
  • Event/push триггер первым; deferrable sensor как сеть безопасности от известного времени релиза. Никаких грубых cron-ов.
  • Жёсткий SLA + эскалация измеряются от времени релиза; промах -> пейджер с runbook (перезапрос, проверить статус вендора).
  • Pre-warm DAG/воркеров; зарезервировать pool, чтобы шумный сосед не заморил голодом критический фид.
  • Быстрый путь валидации: must-have проверки (схема, sanity одной строки, диапазон) на критическом пути; отложить тяжёлые проверки.
  • Fail visible, not silent: для клиента на терминале поздно-но-правильно бьёт быстро-но-неправильно.
hook"Mirrors the McMakler Salesforce sync via Airflow where freshness was business-critical: event trigger + deferrable sensor + dedicated pool + SLA-with-escalation."
якорь«Зеркалит McMakler Salesforce sync через Airflow, где свежесть была бизнес-критична: event триггер + deferrable sensor + выделенный pool + SLA-с-эскалацией.»

7 · Partitioning · Incremental Loads · dbt7 · Партиции · Инкрементальные загрузки · dbt

Partitioning for time seriesПартиции для временных рядов

  • Partition by the dim you filter/reprocess on: usually observation_date (raw often by load_date).
  • Enables cheap incremental + idempotent partition overwrite + scan pruning for as-of queries.
  • Avoid tiny partitions (small-files) and skewed giants.
  • Keep revisions inside the observation_date partition -> a late revision rewrites exactly one partition.
  • Cluster/sort within partition by series_id.
  • Партиция по измерению, по которому фильтруешь/reprocess: обычно observation_date (raw часто по load_date).
  • Разрешает дешёвый инкрементальный + идемпотентный partition overwrite + scan pruning для as-of запросов.
  • Избегать крошечных партиций (small-files) и перекошенных гигантов.
  • Хранить ревизии внутри партиции observation_date -> поздняя ревизия перезаписывает ровно одну партицию.
  • Cluster/sort внутри партиции по series_id.

Incremental vs full loadИнкрементальная vs полная загрузка

  • Full: reprocess everything each run. Fine for small/fragile data; catches drift.
  • Incremental: only new/changed rows -> cheaper at scale; needs correct change key + lookback.
  • Keep a periodic full-refresh to reconcile incremental drift.
  • Полная: reprocess всё каждый запуск. Подходит для маленьких/хрупких данных; ловит drift.
  • Инкрементальная: только новые/изменённые строки -> дешевле на масштабе; нужен правильный change key + lookback.
  • Держать периодический full-refresh, чтобы reconcile инкрементальный drift.

dbt incremental modelsdbt инкрементальные модели

  • Transform only new/changed rows via is_incremental() + a filter.
  • Strategies: append · merge (upsert on unique_key) · delete+insert · insert_overwrite (partition replace).
  • Трансформировать только новые/изменённые строки через is_incremental() + фильтр.
  • Стратегии: append · merge (upsert по unique_key) · delete+insert · insert_overwrite (partition replace).
select series_id, observation_date, value, revision_id, load_ts
from {{ source('vendor','cpi') }}
{% if is_incremental() %}
  -- lookback window catches late revisions to old dates
  where load_ts > (select max(load_ts) from {{ this }}) - interval 3 day
{% endif %}
gotchaToo-tight an incremental filter misses late-arriving rows. Always add a lookback window. Wrong unique_key -> dupes.
gotchaСлишком жёсткий инкрементальный фильтр пропускает поздно прибывающие строки. Всегда добавлять lookback window. Неправильный unique_key -> дубли.

dbt snapshots (SCD-2)dbt snapshots (SCD-2)

  • Capture how a mutable source row changes over time: dbt_valid_from / dbt_valid_to.
  • Gives history / as-of for changing records.
  • Fit: slowly-changing reference data (series definitions, vendor mappings); capturing vintages when a source overwrites in place.
  • Incrementals = efficient fact appends; snapshots = history of change. Different jobs.
  • Захватывает, как мутабельная строка источника меняется со временем: dbt_valid_from / dbt_valid_to.
  • Даёт историю / as-of для меняющихся записей.
  • Подходит для: медленно меняющиеся справочные данные (определения серий, маппинги вендоров); захват vintages, когда источник перезаписывает in place.
  • Инкрементальные = эффективные fact appends; snapshots = история изменений. Разные задачи.
hook"At Career.io I institutionalised dbt code-review + testing standards so quality was a property of the workflow, not a bolt-on."
якорь«В Career.io я институционализировал dbt code-review + стандарты тестирования, чтобы качество было свойством воркфлоу, а не bolt-on.»

8 · Vendor Feeds · Schema Evolution · Monitoring · Correctness8 · Фиды вендоров · Эволюция схем · Мониторинг · Корректность

Heterogeneous vendor cadencesГетерогенные каденции вендоров

  • Feed registry as config: cadence, format, contract, SLA, transport, PK per vendor -> generates the DAG.
  • Normalize early: land raw immutably, conform every vendor to the one canonical model.
  • Conflicting values: explicit precedence/tie-break rule + provenance, so you can explain which source won.
  • Feed registry как конфиг: каденция, формат, контракт, SLA, транспорт, PK на вендора -> генерирует DAG.
  • Нормализовать рано: посадить raw неизменно, привести каждого вендора к одной канонической модели.
  • Конфликтующие значения: явное правило приоритета/tie-break + provenance, чтобы можно было объяснить, какой источник выиграл.
hook"Conforming heterogeneous sources to one model = the Meta 10-system consolidation: mapped overlapping fields to a single schema, cut 47% redundant volume."
якорь«Приведение гетерогенных источников к одной модели = Meta консолидация 10 систем: смаппил перекрывающиеся поля в единую схему, срезал 47% избыточного объёма.»

Schema drift (vendor changes format)Drift схемы (вендор меняет формат)

  • Contract check at the gate: compare incoming schema to expected; drift -> fail/quarantine, don't publish.
  • Tolerate additive (new column) by config; block breaking (rename/remove/retype) until reviewed.
  • Alert + runbook; version the contract in source control; record schema changes in metadata.
  • Проверка контракта на воротах: сравнить входящую схему с ожидаемой; drift -> fail/карантин, не публиковать.
  • Терпеть additive (новая колонка) по конфигу; блокировать breaking (rename/remove/retype) до ревью.
  • Алерт + runbook; версионировать контракт в source control; записывать изменения схемы в метаданные.

Monitoring hundreds of feedsМониторинг сотен фидов

  • Golden data signals: freshness · completeness · volume · validity · schema.
  • Per-feed SLA dashboard + lineage for blast-radius.
  • Beat alert fatigue: alert on breaches not on every run; dedupe, group, severity-tier, route by owner, auto-resolve on next pass.
  • Remediation built in: alerts link a runbook; re-pull / re-run partition is one command.
  • Золотые сигналы данных: свежесть · полнота · объём · валидность · схема.
  • SLA дашборд на фид + lineage для blast-radius.
  • Победить alert fatigue: алертить на нарушения, а не на каждый запуск; dedup, группировать, severity-tier, роутить по владельцу, auto-resolve на следующем проходе.
  • Remediation встроена: алерты ссылаются на runbook; re-pull / перезапуск партиции — одна команда.

Lineage & the wrong-number incidentLineage и инцидент с неправильным числом

  • Lineage at model layer (dbt graph) + pipeline (asset graph / OpenLineage) + load metadata (load_ts, file hash, revision).
  • Client sees wrong CPI: reproduce as-of -> trace back marts -> staging -> raw vintage -> vendor; classify (vendor value / parse-unit bug / join fanout / revision handling); contain (idempotent partition reprocess); prevent (add check + regression test + tighten gate).
  • Lineage на уровне модели (dbt граф) + пайплайна (asset граф / OpenLineage) + метаданные загрузки (load_ts, хеш файла, ревизия).
  • Клиент видит неправильный CPI: воспроизвести as-of -> трассировать назад marts -> staging -> raw vintage -> вендор; классифицировать (значение вендора / баг parse-unit / join fanout / обработка ревизии); изолировать (идемпотентный partition reprocess); предотвратить (добавить проверку + регрессионный тест + затянуть gate).
hook"Same muscle as the Meta entity-resolution fix: validate against what consumers see (100+ metrics), not just the source, then make the fix safe to reprocess."
якорь«Тот же навык, что и Meta фикс entity-resolution: валидировать против того, что видят консьюмеры (100+ метрик), а не только источник, затем сделать фикс безопасным для reprocess.»

Point-in-time correctness (the domain differentiator)Point-in-time корректность (различие домена)

  • Reconstruct what was known at a past moment, before later revisions.
  • Prevents look-ahead bias in backtests / forecast evaluation / client analysis.
  • Implementation: retain all vintages keyed by (series_id, observation_date, publish_ts); never overwrite; expose as-of views.
  • Восстановить, что было известно в прошлый момент, до последующих ревизий.
  • Предотвращает look-ahead bias в backtests / оценке прогнозов / клиентском анализе.
  • Реализация: хранить все vintage с ключом (series_id, observation_date, publish_ts); никогда не перезаписывать; выставлять as-of виды.
hook"At Front Tier on bank/financial data I learned point-in-time discipline: a model trained on revised-after-the-fact numbers is silently cheating. Same applies to macro vintages."
якорь«В Front Tier на банковских/финансовых данных я выучил point-in-time дисциплину: модель, обученная на пересмотренных-после-факта числах, тихо читерит. То же применимо к макро vintages.»
Safe AI in the pipeline: anomaly detection on incoming values; schema inference for messy feeds; LLM-assisted enrichment/validation-rule suggestion. Always human-in-the-loop and guardrailed: AI flags, deterministic gates decide. (My Untangler multi-agent project won Meta's EMEA DE Hackathon and was adopted internally.)
Безопасный AI в пайплайне: детект аномалий на входящих значениях; вывод схемы для грязных фидов; LLM-ассистированное обогащение/предложение правил валидации. Всегда human-in-the-loop и с ограждениями: AI флагает, детерминистичные gates решают. (Мой проект Untangler с multi-agent выиграл Meta EMEA DE Hackathon и был принят внутренне.)

9 · Self-Quiz — Flip to Reveal9 · Самопроверка — переверни карту

Click a card to flip. Answer out loud first.

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

What makes a task idempotent, and why is it the precondition for backfills?
Что делает задачу идемпотентной, и почему это условие для backfills?
tap to reveal
кликни, чтобы увидеть
Same input -> same end state, no dupes/drift. Achieve via partition overwrite / MERGE on a stable key + passing the logical date in. Backfills re-run N partitions in any order/parallel; without idempotency you double-count.
Тот же вход -> то же конечное состояние, никаких дублей/drift. Достигается partition overwrite / MERGE по стабильному ключу + передачей logical date. Backfills перезапускают N партиций в любом порядке/параллельно; без идемпотентности удваивается счёт.
poke vs reschedule vs deferrable sensors?
poke vs reschedule vs deferrable sensors?
tap to reveal
кликни, чтобы увидеть
poke holds a worker slot the whole wait (can starve the pool). reschedule frees the slot between checks (use for minutes+). Deferrable = async via triggerer, zero slot -> best for "wait for the CPI file". Always set a timeout.
poke держит worker slot всё ожидание (может заморить pool голодом). reschedule освобождает slot между проверками (для ожидания минут+). Deferrable = async через triggerer, нулевой slot -> лучше всего для «ждать файл CPI». Всегда ставить timeout.
How do you model revisions so as-of queries work?
Как моделировать ревизии, чтобы as-of запросы работали?
tap to reveal
кликни, чтобы увидеть
Grain = (series_id, observation_date, revision/vintage, publish_ts). Append new vintages, never overwrite. Current = latest revision per (series,date). As-of = latest revision with publish_ts <= X. Avoids look-ahead bias.
Grain = (series_id, observation_date, revision/vintage, publish_ts). Append новые vintages, никогда не перезаписывать. Текущее = последняя ревизия на (series,date). As-of = последняя ревизия с publish_ts <= X. Избегает look-ahead bias.
Biggest failure mode of dbt incremental models?
Главный режим отказа dbt инкрементальных моделей?
tap to reveal
кликни, чтобы увидеть
Too-tight incremental filter misses late-arriving rows -> add a lookback window. Also: wrong unique_key -> dupes; logic drift vs full-refresh -> periodically full-refresh to reconcile.
Слишком жёсткий инкрементальный фильтр пропускает поздно прибывающие строки -> добавить lookback window. Также: неправильный unique_key -> дубли; drift логики vs full-refresh -> периодически full-refresh для reconcile.
Snapshot vs incremental model — when each?
Snapshot vs incremental модель — когда каждая?
tap to reveal
кликни, чтобы увидеть
Snapshot = SCD-2 history of a mutable record (valid_from/valid_to) -> reference/metadata, vintages. Incremental = efficient append/merge of new facts. Different jobs.
Snapshot = SCD-2 история мутабельной записи (valid_from/valid_to) -> справочники/метаданные, vintages. Incremental = эффективный append/merge новых фактов. Разные задачи.
How do you guarantee CPI lands minutes after the print?
Как гарантировать, что CPI приземлится через минуты после публикации?
tap to reveal
кликни, чтобы увидеть
Release calendar as data; event/push trigger + deferrable sensor safety-net; tight SLA + escalation from release time; pre-warm + dedicated pool; fast critical-path validation; fail visible. Late-but-correct beats fast-but-wrong.
Календарь релизов как данные; event/push триггер + deferrable sensor как сеть; жёсткий SLA + эскалация от времени релиза; pre-warm + выделенный pool; быстрая critical-path валидация; fail visible. Поздно-но-правильно бьёт быстро-но-неправильно.
When pick Temporal over Airflow/Dagster?
Когда выбирать Temporal вместо Airflow/Dagster?
tap to reveal
кликни, чтобы увидеть
Temporal = durable execution for long-running, stateful, fault-tolerant workflows-as-code (e.g. multi-step human-in-the-loop onboarding). Not a data scheduler. For hundreds of scheduled validated batch feeds use Airflow (or Dagster for asset+lineage+DQ native).
Temporal = durable execution для долгих, stateful, fault-tolerant workflows-as-code (напр. multi-step human-in-the-loop онбординг). Не дата-шедулер. Для сотен запланированных валидированных батч-фидов использовать Airflow (или Dagster для asset+lineage+DQ native).
Client sees a wrong number. First moves?
Клиент видит неправильное число. Первые шаги?
tap to reveal
кликни, чтобы увидеть
Reproduce as-of -> trace lineage back marts->staging->raw vintage->vendor -> classify (vendor / parse-unit / join fanout / revision) -> contain (idempotent partition reprocess) -> prevent (add check + regression test + tighten gate).
Воспроизвести as-of -> трассировать lineage назад marts->staging->raw vintage->vendor -> классифицировать (vendor / parse-unit / join fanout / revision) -> изолировать (idempotent partition reprocess) -> предотвратить (добавить проверку + регрессионный тест + затянуть gate).
Which Airflow executor for bursty, isolated workloads & how to protect a vendor rate limit?
Какой Airflow executor для bursty, изолированных нагрузок и как защитить rate limit вендора?
tap to reveal
кликни, чтобы увидеть
Kubernetes executor (1 pod/task, isolation, bursty) — accept pod startup latency. Protect a rate limit with a Pool sized to the limit; tune parallelism / max_active_tasks; push heavy compute off the worker.
Kubernetes executor (1 pod/задача, изоляция, bursty) — принять задержку старта pod. Защитить rate limit через Pool размером в лимит; тюнить parallelism / max_active_tasks; пушить тяжёлые вычисления с воркера.
Modernize a fragile legacy pipeline without breaking consumers?
Модернизировать хрупкий legacy пайплайн, не сломав консьюмеров?
tap to reveal
кликни, чтобы увидеть
Map consumers first -> strangler-fig new path alongside old -> parallel run -> reconcile (counts/key-diffs/tolerances) -> cut over consumer-by-consumer with fallback -> decommission. Embed tests/version control so debt doesn't return. (= McMakler triple migration.)
Смаппить консьюмеров первыми -> strangler-fig новый путь рядом со старым -> параллельный запуск -> reconcile (counts/key-diffs/tolerances) -> переключать консьюмера за консьюмером с fallback -> вывести из эксплуатации. Встроить тесты/version control, чтобы долг не вернулся. (= McMakler тройная миграция.)