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
- 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 duckdbstg_klines/stg_funding, thenfeatures_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:
begin;delete from stg_klines where date = $day; -- the partition, not the row, is the unitinsert 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):
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
- 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.
- partition by event-time day (+ entity), deterministic from the data. never partition by run date or
now(). - replace partitions atomically. files: write to temp, then
os.replace. tables: delete+insert in one transaction, orMERGE/ dltwrite_disposition="merge"with aprimary_key(dlt's merge-loading guide). also declare a primary key so an accidental double insert fails loudly. - keep raw as text/int, and decide units and types in staging. a schema change then breaks one sql expression, not the landing.
- record lineage in the data:
source_sha256on raw rows makes "which partitions changed at the source?" a query (the late-data case). - 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. - 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.
- 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). - retries: exponential backoff + full jitter, honour
Retry-After, retry only 418/429/5xx/timeouts. a 400 is a bug; retrying it only hides it. - 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
- binance spot kline timestamps changed from ms to µs on 2025-01-01 (spot only; usdt-m futures files stayed ms). a naive
ms*1000gives year 56971. detect by magnitude (binance-public-data#386). - spot csvs have no header; futures csvs do (
open_time,open,...).int('open_time')raises. detect the header, don't assume it. - funding archive columns ≠ rest fields. the archive has
calc_time,funding_interval_hours,last_funding_rate;/fapi/v1/fundingRatereturnsfundingTime,fundingRate,markPrice,rateType. treat them as two sources with two contracts. - funding timestamps jitter by milliseconds (
...600015). exact-time joins silently miss. use asof. - 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.
- 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.
- 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. - websocket latency measured with your clock lies:
recv - event_timemedian was -183 ms (the local clock was behind binance's). compare event time only against event time. - airflow 3
catchupdefaults 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). - duckdb column aliases without
ASbreak on keywords (close,rows,days), e.g.x::double closefails to parse. always writeAS.
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_BYwithOVERWRITE_OR_IGNOREoverwrites 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_datemarks runs success without executing. - dagster-specific bugs: the dagster page's known bugs.
troubleshooting
| symptom | cause | fix |
|---|---|---|
| timestamps in year 56971 from 2025 data | spot open_time switched to µs | normalize by magnitude in staging; add a "ts within partition day" contract check |
| ValueError: invalid literal for int() ... 'open_time' | futures csv has a header row | detect a non-numeric first cell and skip it |
| row count 2x/3x after a rerun, volume doubled | append instead of replace | delete+insert per partition in one txn, plus a primary key |
| a corrected day is fixed but next-day features stay old | lookback dependency not propagated | declare partition mapping (offset -1); rerun d..d+window |
| gap check fails but reconciliation passes | both compare against the same (gappy) source | keep the completeness check (expected 1440/day) separate from transport reconciliation |
| IO Error: Could not set lock on file ... warehouse.duckdb | two processes open the duckdb file for write | in-process executor or one writer; read-only connections for readers |
| crash left some symbols written, others not | multi-entity partition written in a loop | atomic 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 keyword | expr 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 nullrv_24hvalues 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-rangerefused because the assets lackBackfillPolicy.single_run(), so the backfill ran as a per-partition loop (seedagster). - (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 changedstg_klinesandfeatures_1hfor 01-03 only, and 01-02 stayed byte-identical. 01-04features_1hstayed 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
1735603200000and for 2025-01-01 raw1735689600000000. naivems*1000andepoch_ms()both gave56971-10-25, while normalized staging gave2025-01-01 00:00:00. a futures header row raised aValueErroruntil it was detected. the contract checkwrong_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_completefailed (blocking),stg_klines_reconcilespassed (same-source reconciliation cannot see source gaps), andfeatures_1hnever 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_1mproduced 59 events. after a simulated replay doubled the log (118 events), dedup on(symbol, open_time), preferringclosed 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.connectfor write reproducedCould not set lock on file, so the example pinsin_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
| need | default | alternatives / when |
|---|---|---|
| orchestration | dagster 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 apis | dlt (its rest_api, incremental, and schema-contract guides) | hand-rolled requests when the api is one endpoint (this case); airbyte/fivetran for saas sources |
| transform | sql in duckdb, or dbt (dbt) | polars for python-shaped logic |
| storage | duckdb file + hive parquet (olap) | clickhouse for real-time/high concurrency, postgres for oltp-shaped serving (postgres), iceberg/delta/ducklake for multi-engine lakes |
| streaming | none (append log + micro-batch) | kafka/redpanda once there are several consumers, replay needs, or sub-minute slas (the kafka design guide goes deeper) |
| quality | orchestrator checks (sql) + dbt tests | pandera (dataframes), great expectations / soda (contracts as config) |
| cdc | debezium -> kafka | only when the source is an oltp db you own |
| features / ml | feature-engineering, ml | domain reasoning for crypto data: crypto |