data-pipeline trained

versions practiced: python 3.12, dagster 1.13.24, duckdb 1.5.5, pyarrow 25.0.1, dbt-core 1.12.5 + dbt-duckdb 1.11.0, dagster-dbt 0.29.24, websockets 17.1. orchestrator depth lives in dagster; this page stays orchestrator-neutral.

mental model

source
landingexact bytes, checksum
rawtyped-as-text, partitioned
stagedtyped, normalized, deduped
quality gatesblock
curated / features
consumers
the partition (for example symbol by utc day) is the unit of work: every step replaces its partition, and the gates read raw and staged before anything is promoted.
  • a pipeline is a pure function per partition. output(partition) = f(inputs(partition + declared lookback), code). rerun = overwrite, never append. this is beauchemin's "functional data engineering": the partition is the atomic unit of state and immutable raw data lets you rebuild anything (beauchemin).
  • exactly-once is a property of the sink, not the transport. delivery is at-least-once in practice; "exactly-once" = at-least-once delivery + idempotent writes keyed by a natural id (confluent eos post).
  • event time decides the partition; processing time decides nothing. a candle belongs to the day of its open_time, whenever it arrived.
  • every partition has a dependency window. a 24h rolling feature for day d reads d-1. so a correction to d-1 makes d stale too. declare the window; propagate reruns through it.
  • warehouse is enough until one of: many writers, many engines, or data > one machine. duckdb + hive parquet covered the whole practice case; routing for engines and table formats is in olap.

examples

crypto-market-pipeline is verified working end to end:

  • pipeline.py: binance spot 1m daily zips (data.binance.vision, checksum-verified) and usdt-m funding over rest (cursor pagination, backoff). output: raw hive parquet, then duckdb stg_klines / stg_funding, then features_1h. orchestrated as dagster daily partitions with blocking checks.
  • scenarios.py: the six correctness proofs below.
  • ws_ingest.py: a live websocket append log, deduplicated on replay.

the partition-replace core, the one pattern to copy:

partition replace
begin;
delete from stg_klines where date = $day; -- the partition, not the row, is the unit
insert into stg_klines select ... from read_parquet($raw_glob_for_day);
commit; -- readers never see a half partition

timestamp unit normalization (binance spot moved ms -> µs on 2025-01-01):

unit normalization
case when open_time >= 100000000000000 then make_timestamp(open_time) -- µs (16 digits)
else make_timestamp(open_time * 1000) end -- ms (13 digits)

best practices

  1. land the exact source bytes first (zip + published checksum). raw is replayable forever. you can reprocess without re-hitting the api, and you can prove what the source said.
  2. partition by event-time day (+ entity), deterministic from the data. never partition by run date or now().
  3. replace partitions atomically. files: write to temp, then os.replace. tables: delete+insert in one transaction, or MERGE / dlt write_disposition="merge" with a primary_key (dlt's merge-loading guide). also declare a primary key so an accidental double insert fails loudly.
  4. keep raw as text/int, and decide units and types in staging. a schema change then breaks one sql expression, not the landing.
  5. record lineage in the data: source_sha256 on raw rows makes "which partitions changed at the source?" a query (the late-data case).
  6. declare lookback windows (TimeWindowPartitionMapping(start_offset=-1) in dagster) so the orchestrator knows d depends on d-1, and rerun downstream partitions after a correction.
  7. gates at every boundary, and blocking: completeness (1440 min/day), gaps, uniqueness, reconciliation vs source, freshness. a failed gate must stop downstream in the same run, not just raise an alert.
  8. point-in-time joins for features: ASOF JOIN ... on bar_close >= event_time. never use equality on event timestamps. funding times jitter (1735689600015) (feature leakage rules: feature-engineering).
  9. retries: exponential backoff + full jitter, honour Retry-After, retry only 418/429/5xx/timeouts. a 400 is a bug; retrying it only hides it.
  10. choose the orchestrator by abstraction: asset/partition-centric (dagster) for data products. task-centric (airflow) when the org already runs it. python-flow-centric (prefect) for app-like workflows. details below and in dagster.

strengths

  • a partitioned, idempotent batch design makes backfills, corrections, and crash recovery the same operation: rerun the partition. all of these were proven with hashes below.
  • duckdb + parquet handles tens of millions of rows on a laptop with zero infra. this includes asof joins, window features, and reading hive globs.
  • public crypto archives (data.binance.vision) publish per-file sha256 checksums, so ingestion can be verified, not trusted.

weaknesses / pain points

  • lookback windows make partitions non-independent. a correction ripples forward by the window length. most orchestrators do not auto-rerun that ripple unless you declare it and turn on automation.
  • reconciling against the same source cannot catch source gaps. it proves transport, not completeness (the quality-gate case below).
  • duckdb is single-writer per file. parallel steps in separate processes contend for the lock. serialize writes or move to a server warehouse or a lakehouse catalog (olap).
  • streaming adds cost everywhere: offsets, state, replays, clock skew, operational load. for daily/hourly ml features, micro-batch over an append log wins.

gotchas

  1. binance spot kline timestamps changed from ms to µs on 2025-01-01 (spot only; usdt-m futures files stayed ms). a naive ms*1000 gives year 56971. detect by magnitude (binance-public-data#386).
  2. spot csvs have no header; futures csvs do (open_time,open,...). int('open_time') raises. detect the header, don't assume it.
  3. funding archive columns ≠ rest fields. the archive has calc_time,funding_interval_hours,last_funding_rate; /fapi/v1/fundingRate returns fundingTime,fundingRate,markPrice,rateType. treat them as two sources with two contracts.
  4. funding timestamps jitter by milliseconds (...600015). exact-time joins silently miss. use asof.
  5. append-mode reruns double silently: 3 runs gave 12,960 rows and 3x volume, with 4,320 distinct keys (the late-data case below). a declared primary key turns this into a loud error.
  6. a correction to day d leaves d+1 stale when d+1 has a lookback. the d+1 hash only changed after its own rerun.
  7. a retry policy on a deterministic error just burns time. RetryPolicy(max_retries=2) retried a simulated crash twice ("exceeded max_retries of 2"). classify errors before retrying.
  8. websocket latency measured with your clock lies: recv - event_time median was -183 ms (the local clock was behind binance's). compare event time only against event time.
  9. airflow 3 catchup defaults to false (scheduler.catchup_by_default=False), not the old true. older blog advice to "set catchup=false" is now redundant, and missing backfills surprise people the other way (the airflow docs cover this under dag runs).
  10. duckdb column aliases without AS break on keywords (close, rows, days), e.g. x::double close fails to parse. always write AS.

known bugs

the upstream issues that bit or shaped this design:

  • binance-public-data#386: the µs switch. #474: an apparent "volume cliff" that was really a misread column or unit.
  • duckdb#10282: COPY ... PARTITION_BY with OVERWRITE_OR_IGNORE overwrites files, not the partition directory. stale files can survive. delete the partition dir or write per-partition files yourself (as the example does).
  • airflow#38513: a future start_date marks runs success without executing.
  • dagster-specific bugs: the dagster page's known bugs.

troubleshooting

symptomcausefix
timestamps in year 56971 from 2025 dataspot open_time switched to µsnormalize by magnitude in staging; add a "ts within partition day" contract check
ValueError: invalid literal for int() ... 'open_time'futures csv has a header rowdetect a non-numeric first cell and skip it
row count 2x/3x after a rerun, volume doubledappend instead of replacedelete+insert per partition in one txn, plus a primary key
a corrected day is fixed but next-day features stay oldlookback dependency not propagateddeclare partition mapping (offset -1); rerun d..d+window
gap check fails but reconciliation passesboth compare against the same (gappy) sourcekeep the completeness check (expected 1440/day) separate from transport reconciliation
IO Error: Could not set lock on file ... warehouse.duckdbtwo processes open the duckdb file for writein-process executor or one writer; read-only connections for readers
crash left some symbols written, others notmulti-entity partition written in a loopatomic per-file writes make the rerun converge; the rerun rewrote all 3 symbols, 1440 rows each
Parser Error: syntax error at or near "close"alias without AS on a duckdb keywordexpr AS close

practiced cases

on macos, data from data.binance.vision + fapi.binance.com (public, no keys), btcusdt/ethusdt/solusdt, 2024-12-29..2025-01-04 (crossing the µs switch). example: crypto-market-pipeline.

  • (a) end-to-end build. 21 zips (about 70 kb each), 7 daily partitions via dagster asset materialize --partition, 7/7 runs succeeded (about 13 s each, mostly cli start plus polite funding pagination). row counts:
    • stg_klines: 30,240 rows (7 x 3 x 1440).
    • stg_funding: 63 rows (3 per day per symbol, fetched at page size 2, so 2 pages each).
    • features_1h: 504 rows. the only null rv_24h values are the 6 warm-up rows of the first partition, and 0 funding values are null.
  • (b) idempotency. 2025-01-01 materialized 3 times: raw parquet file sha256 and table partition hashes were identical (stg_klines ab4dba2e7ffeff4b, features_1h 8512b93ab8f21ce5). crash/resume: simulated failure before ethusdt on 2025-01-02. the run failed (after 2 wasted retries on the simulated crash) with btc/sol raw written, eth missing, and downstream not run. a plain rerun succeeded with 1440/1440/1440 rows and all checks passed. backfill: --partition-range refused because the assets lack BackfillPolicy.single_run(), so the backfill ran as a per-partition loop (see dagster).
  • (c) late / corrected data. rewrote solusdt 2025-01-03 (one minute's volume +1000, last close +0.5%). the sha-diff "sensor" flagged exactly ['2025-01-03']. rerunning it changed stg_klines and features_1h for 01-03 only, and 01-02 stayed byte-identical. 01-04 features_1h stayed stale until its own rerun (then its hash changed). the run repeated twice with the same outcome. without idempotent overwrite: append x3 gave 4,320, then 8,640, then 12,960 rows, with 4,320 distinct keys and volume 1,855,772, then 3,711,544, then 5,567,316.
  • (d) schema drift. for 2024-12-31 raw 1735603200000 and for 2025-01-01 raw 1735689600000000. naive ms*1000 and epoch_ms() both gave 56971-10-25, while normalized staging gave 2025-01-01 00:00:00. a futures header row raised a ValueError until it was detected. the contract check wrong_day_ts (min/max ts inside the partition day) is the gate that catches a unit error.
  • (e) quality gates. deleted 5 minutes from btcusdt 2025-01-04 at the source. stg_klines_complete failed (blocking), stg_klines_reconciles passed (same-source reconciliation cannot see source gaps), and features_1h never started. after restoring the source, the rerun succeeded. the freshness check ran as a warn-severity check.
  • (f) live websocket. 65 s on btcusdt@kline_1m + ethusdt@kline_1m produced 59 events. after a simulated replay doubled the log (118 events), dedup on (symbol, open_time), preferring closed desc, event_time desc, gave 4 klines (2 closed, 2 open). clock skew was a median of -183 ms, as the latency gotcha above predicts.
  • multiprocess executor vs duckdb. in one run the default multiprocess executor succeeded but spent about 2–6 s spawning each step subprocess (about 80 s for one partition, against about 13 s in-process). a second concurrent duckdb.connect for write reproduced Could not set lock on file, so the example pins in_process_executor.
  • not practiced (unverified here): dlt as the loader, pandera/great expectations/soda, kafka/redpanda, iceberg/delta/ducklake, debezium cdc, scd2 snapshots in practice (dbt snapshots were practiced in dbt), airflow and prefect runtimes.

ecosystem

needdefaultalternatives / when
orchestrationdagster assets + partitions (dagster)airflow when the org runs it, and airflow 3 now has assets (the airflow docs cover them). prefect for python-first flows with task caching (the prefect docs cover caching). cron + make for one job.
declarative loading from apisdlt (its rest_api, incremental, and schema-contract guides)hand-rolled requests when the api is one endpoint (this case); airbyte/fivetran for saas sources
transformsql in duckdb, or dbt (dbt)polars for python-shaped logic
storageduckdb file + hive parquet (olap)clickhouse for real-time/high concurrency, postgres for oltp-shaped serving (postgres), iceberg/delta/ducklake for multi-engine lakes
streamingnone (append log + micro-batch)kafka/redpanda once there are several consumers, replay needs, or sub-minute slas (the kafka design guide goes deeper)
qualityorchestrator checks (sql) + dbt testspandera (dataframes), great expectations / soda (contracts as config)
cdcdebezium -> kafkaonly when the source is an oltp db you own
features / mlfeature-engineering, mldomain reasoning for crypto data: crypto

search pages

go to any page