Data Pipeline Observatory
A production-style ETL pipeline on 1.07 million real invoice lines from a UK online gift retailer. It validates every row, quarantines what it cannot trust with a reason, loads PostgreSQL without duplicates and refuses to publish a batch that fails its quality gate.
1Problem
Most companies run reports on data that arrives from several systems that do not agree. Here, a retailer has three: a monthly invoice export, a product catalogue and a CRM with its customer register. The export is real, and so are its problems:
- Cancellations recorded as separate invoices with negative quantities.
- Exact duplicate lines, including thousands from overlapping export periods.
- Stock write-offs (damaged, missing, re-labelled) recorded as zero-price lines with no customer.
- Bad-debt adjustments, test records, postage and manual lines mixed in with product sales.
- More than a fifth of lines with no customer ID, and stock codes in inconsistent case.
Summed naively, this data overstates sales, double-counts customers and misattributes revenue. Nobody notices until a number in a board pack is wrong.
2Why it matters
Analysts typically spend more time cleaning and reconciling data than analysing it, and the cleaning is usually manual, undocumented and repeated every month. When it happens in a spreadsheet, nobody can say which rows were removed or why.
A pipeline that applies the same written rules every time, keeps every rejected row with its reason, and stops rather than publishing bad data turns that monthly chore into something auditable. It is also the foundation the later projects on this site build on: the CSV analyst, the document assistant and the support agent all read this warehouse.
3Solution
Each month's file goes through five steps. Rows are checked against a data contract and a set of business rules; every rule has a reason code and a plain explanation. Failing rows go to quarantine with their original content and line number. Rows that are real but imperfect (a guest sale with no customer ID) are published with a warning instead of being thrown away.
Before anything is written, batch-level checks score the drop. A failed gate check stops the run. Passing drops are loaded in one transaction, reconciled against the warehouse, and recorded with their steps, log lines and check results so any past run can be inspected.
4Architecture
Monthly CSV drops
25 files, original export layout
Product catalogue
JSON master data
CRM REST API
paginated, can fail
1Extract
Retries with backoff and jitter; structural errors caught per line
2Validate
Data contract, business and integrity rules with reason codes, de-duplication
3Transform
Normalise codes and text, classify lines, stable line keys
4Quality gate
Batch checks; a failed gate stops the run before any write
5Load
COPY + upserts in one transaction, reconciliation check
Quarantine
Every rejected row kept with its line number and reason code
PostgreSQL warehouse
retail.invoice, invoice_line, customer, product; read by later projects
Run metadata
Steps, real log lines, data-quality results, environment
Triggers. A schedule every 6 hours, visitors (rate-limited), or the CLI.
Execution. A PostgreSQL job queue and a global CPU lease, so one heavy job runs at a time on 2 vCPUs.
Read side. A FastAPI service feeds the dashboard below; nginx in front limits request rates.
The product catalogue and CRM register are derived from the same public dataset but served as separate sources, a JSON file and a paginated REST API, so the pipeline faces the integration problems of a real setup: joins across systems, pagination and an upstream service that can fail.
5Technical implementation
- Pipeline
- Python 3.13, pandas, tenacity, a data contract and reason-code registry in plain Python
- Storage
- PostgreSQL 16: warehouse schema, quarantine, run metadata and the job queue
- Jobs
- Procrastinate (PostgreSQL-backed queue) with a periodic schedule and a global CPU lease
- API
- FastAPI and Pydantic, typed responses, RFC 9457 error bodies, per-visitor quotas
- Interface
- Next.js with types generated from the API's OpenAPI document
- Operations
- Docker images built and scanned in CI, nginx, structured JSON logs, Prometheus metrics
Data contract. The header must match an 8-column contract exactly. Lines are parsed one at a time, so a single truncated line or invalid byte sequence becomes one quarantined row rather than a failed file.
Rules with reason codes. 18 rejection rules and 5 warning rules, each with a code, a severity and an explanation shown in the dashboard. Rules are evaluated in a fixed order, so every quarantined row has exactly one primary reason.
Idempotent loads. Each line's key is a SHA-256 hash of its normalised content. Rows are bulk-copied into temporary tables and inserted with ON CONFLICT DO NOTHING, so reprocessing never creates duplicates.
Quality score. 9 checks across validity, uniqueness, completeness, timeliness, integrity and accuracy. Score = 100 × Σ(weight × passed) / Σ(weight), with weights of 3 for gate checks, 2 for warnings and 1 for information. The formula is shown under every run.
6Live demonstration
Start a run and watch it move through the steps. The standard scenario loads the next month of real invoices; the simulated scenarios inject failures so you can see how each one is handled. Select any run to see its checks, quarantined rows and log.
Pipeline dashboard
Portfolio demonstration on public data. Figures are read live from the running system.
Last run
No runs yet
All runs
No successful run yet
- Failed, of the last 0 runs
- 0
- Includes simulated faults
- Monthly drops loaded, of 25
- 0
- Invoice lines in warehouse
- 0
- Customers in warehouse
- 0
- 0 invoices
Loading scenarios…
Recent runs
No runs yet. Start one above, or wait for the next scheduled run.
Start a run to see each step, check and quarantined row here.
7Results
Results are read from the running system and are not available right now. They return when the pipeline service is reachable again.
8Failure handling
- The CRM API is slow or returns errors
- Each page request is tried up to 4 times, with exponential backoff and jitter between attempts. Every retry is logged and counted on the run. Try the Flaky CRM API scenario.
- The CRM is down for the whole run
- Retries are exhausted, the run fails at the extract step and nothing is written. The warehouse keeps the last good state. Try the CRM outage scenario.
- A file arrives damaged
- Damaged lines are caught one at a time (wrong field count, invalid bytes, wrong date format, text in a number column) and quarantined with their line number. If the quarantined share passes 5%, the quality gate stops the run before any write. Try the Corrupted file drop scenario.
- The same file is processed twice
- Every line has a stable key derived from its content, and loads insert only keys that are not already present. Replaying a drop writes 0 new rows, and the run says so.
- A load is interrupted or incomplete
- Each drop loads in one transaction. After the load, a reconciliation check counts the drop's lines in the warehouse and fails the run if any are missing, rolling the transaction back.
- A worker dies mid-run
- A housekeeping job marks runs stuck in the running state for too long as failed, so the dashboard never shows a run that will not finish.
- Too many visitors start runs
- Each visitor can start 3 runs per hour, at most 4 runs can wait in the queue, and nginx limits request rates. Runs execute one at a time under a CPU lease so the rest of the site stays responsive.
9Deployment
The pipeline runs as a worker container next to the API, the website, nginx and PostgreSQL on a single 2 vCPU, 8 GB virtual server. Every push runs linting, type checks and tests, including database tests against a real PostgreSQL. Release images are built once, scanned for vulnerabilities and published with a software bill of materials.
Deployment pulls the tagged images, applies database migrations, restarts the services and checks health through nginx. If the health check fails, the previous version is restored automatically. A scheduled run every 6 hours keeps the dashboard current, and the database is backed up every night.
10Cost
$0 in API spend. The pipeline uses no paid services or AI models. It shares the site's existing server; the worker container is limited to 1 CPU and 1.5 GB of memory, and a month of invoices (30,000 to 80,000 lines) typically processes in a few seconds.
For a real deployment the same design runs on any machine with Docker and PostgreSQL. The main cost driver would be data volume: loads use PostgreSQL's bulk COPY, so they scale with disk throughput rather than per-row round trips.
11Limitations
- The data is historical (December 2009 to December 2011) and is shown with its real dates. It is not shifted to look recent.
- The catalogue and CRM are derived from the invoice data, so they cannot disagree with it in all the ways real systems do. The fault scenarios cover the important failure modes, and are clearly labelled as simulated.
- pandas holds one monthly file in memory. That is appropriate at this volume; much larger drops would call for chunked reads or a columnar engine such as DuckDB or Polars.
- Runs execute one at a time by design, to protect a small shared server. A busy production system would use dedicated workers.
12Source code
The full source is on GitHub under the MIT licence, together with the tests and the infrastructure that deploys it.
- contract.pyData contract and line-by-line parsing
- rules.pyReason-code registry and rule order
- validate.pyBusiness, period and referential rules, de-duplication
- load.pyBulk copy, upserts and reconciliation
- quality.pyData-quality checks and score
- runner.pyRun orchestration and step recording
Tests cover the contract, each rule, idempotent replays, every fault scenario, and a regression test that runs the real December 2010 drop and checks its exact counts.
Data: Chen, D. (2012). Online Retail II [Dataset]. UCI Machine Learning Repository, CC BY 4.0. Details on the data page.