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 стриминг
| Batch | Streaming | |
|---|---|---|
| Data | bounded, periodic | unbounded, continuous |
| Latency | minutes–hours | sub-second–seconds |
| Replay | easy (re-run partition) | harder (offsets, watermarks) |
| Fit | most macro releases | ticks, clickstream, LLM events |
| Батч | Стриминг | |
|---|---|---|
| Данные | ограниченные, периодичные | бесконечные, непрерывные |
| Задержка | минуты–часы | доли секунды–секунды |
| Повтор | легко (перезапуск партиции) | сложнее (оффсеты, watermarks) |
| Подходит для | большинства макрорелизов | тики, кликстрим, события LLM |
The layered pattern everything reduces toСлоёный паттерн, к которому всё сводится
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 фид от вендора → валидированный датасет
Say it as a DAG in this orderРассказывай как DAG в таком порядке
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.
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.
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, затем финализация.
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.
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 списку (одна задача на файл/партицию).
- Подходит для многих гомогенных единиц; не подходит, когда счёт взрывается или каждой нужна кастомная логика.
Executors & scaling workersExecutors и масштабирование воркеров
| Executor | Model | Use when | Watch out |
|---|---|---|---|
| Sequential | one task (SQLite) | local debug | never prod |
| Local | subprocesses, 1 host | small single-node | host is ceiling |
| Celery | distributed workers + broker | classic horizontal scale | broker + fleet ops |
| Kubernetes | 1 pod / task | bursty, isolation | pod startup latency |
| CeleryKubernetes | both | mixed workloads | complexity |
| Executor | Модель | Использовать когда | Следить за |
|---|---|---|---|
| Sequential | одна задача (SQLite) | локальный debug | никогда в prod |
| Local | subprocesses, 1 хост | маленький single-node | хост — потолок |
| Celery | распределённые воркеры + брокер | классический horizontal масштаб | брокер + fleet ops |
| Kubernetes | 1 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
5 · Orchestration Alternatives & Trade-offs5 · Альтернативы оркестрации и trade-offs
| Tool | Core model | Strengths | Weaknesses / fit |
|---|---|---|---|
| Airflow | task DAGs, schedule-centric | mature, huge ecosystem, ubiquitous, strong batch scheduling | task- not asset-centric (historically); streaming/low-latency awkward; ops weight |
| Dagster | software-defined assets, data-aware | assets+lineage first-class, strong typing, great local dev, asset checks (DQ built in) | smaller ecosystem; newer ops patterns |
| Prefect | dynamic Python flows, code-first | very Pythonic, dynamic at runtime, light to start, hybrid execution | less of a batch "platform"; structure can hide in code |
| Temporal | durable execution, workflows-as-code | durable state, retries/timeouts as primitives, long-running stateful + service orchestration | not a data scheduler; you build data semantics yourself |
| Инструмент | Основная модель | Сильные стороны | Слабые стороны / подходит для |
|---|---|---|---|
| Airflow | task DAG-и, schedule-centric | зрелый, огромная экосистема, вездесущий, сильный batch шедулинг | task-, а не asset-centric (исторически); стриминг/low-latency неудобно; операционный вес |
| Dagster | software-defined assets, data-aware | assets+lineage first-class, сильная типизация, отличный local dev, asset checks (DQ встроен) | меньшая экосистема; новые ops паттерны |
| Prefect | динамические Python flows, code-first | очень Pythonic, динамика runtime, лёгкий старт, hybrid исполнение | менее batch «платформа»; структура может спрятаться в коде |
| Temporal | durable execution, workflows-as-code | durable state, retries/timeouts как примитивы, долгоживущие stateful + оркестрация сервисов | не дата-шедулер; data семантика строится вручную |
6 · Scheduling Time-Sensitive Releases6 · Шедулинг релизов с жёсткими SLA
- 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: для клиента на терминале поздно-но-правильно бьёт быстро-но-неправильно.
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 %}
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 = история изменений. Разные задачи.
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, чтобы можно было объяснить, какой источник выиграл.
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).
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 виды.
9 · Self-Quiz — Flip to Reveal9 · Самопроверка — переверни карту
Click a card to flip. Answer out loud first.
Кликни на карту, чтобы перевернуть. Сначала ответь вслух.