data-pipeline
the plumber of truth
a pipeline is a pure function per partition: the same inputs and code give the same output, and a rerun replaces the partition instead of appending to it. land the exact source bytes, derive the partition from event time, and gate every boundary. backfills, corrections, and crash recovery then all become one operation: rerun the partition.
Use when designing, building, reviewing, or debugging a data pipeline -- batch vs micro-batch vs streaming, ELT vs ETL, medallion/layered raw -> staged -> curated, lakehouse vs warehouse, idempotency, deterministic partitioning, exactly-once vs at-least-once + dedup, upserts/merge, late and out-of-order data, event time vs processing time, watermarks, backfills and reprocessing, schema drift and data contracts, CDC (Debezium), slowly changing dimensions, time zones and timestamp units, REST pagination/rate limits/retries, websocket ingest, dlt, Kafka/Redpanda basics, data-quality gates (freshness, row count, gaps, reconciliation), lineage, observability, and orchestrator choice (Dagster vs Airflow vs Prefect). Flagship case -- a crypto market-data pipeline (Binance klines + funding -> Parquet -> DuckDB -> ML features).
methodology
- write the contract: consumers, grain, freshness SLA, and the partition key (event-time day × entity). default to batch. stream only for sub-minute SLAs or several consumers.
- land raw bytes immutably with checksums. keep raw un-coerced. decide types, units, and time zones once, in staging.
- make every write an atomic partition replace with a declared natural key. prove it by running twice and comparing hashes.
- declare lookback windows and late-data handling. detect source changes (sha, CDC,
updated_at) and rerun the changed partition plus its downstream window. - add blocking gates: completeness and gaps, uniqueness, validity, reconciliation vs source, and freshness. fail one on purpose to prove downstream stops.
- orchestrate as assets and partitions (
dagster). transform in SQL/dbt (dbt). store perolap. - review with the design checklist. record lessons in trained.
contents
- trainedlearned layer, crypto practice numbers, gotchas; read first
- design-checklistdesigning or reviewing any pipeline
- crypto-market-pipelinerunnable Binance -> Parquet -> DuckDB -> features pipeline with idempotency, late-data, drift, gate proofs
- dagsterorchestration: assets, partitions, checks, sensors, schedules, dagster-dbt
- dlt-docsofficial dlt: incremental, merge, schema evolution/contracts, rest_api source
- airflow-docsofficial Airflow 3: DAG runs, catchup, backfill, assets, timetables
- prefect-docsofficial Prefect 3: flows, tasks, caching, retries, deployments
- kafka-designofficial Kafka design: logs, consumer offsets, delivery semantics
- known-bugsBinance data, DuckDB, dlt, Airflow issues with workarounds
- sourcesprovenance: repos, commits, canonical writing (Beauchemin, Confluent EOS, DDIA)
- communitycommunity agent skills from Astronomer, dbt Labs, and dlt
- awesome-data-engineeringtool landscape, with awesome-pipeline
related gurus: dbt (transforms), olap (DuckDB, ClickHouse, lakehouse), postgres (OLTP sources and serving), feature-engineering (point-in-time features), ml (training on the output), crypto (market-data meaning).