Skip to content
Uzair Azhar
Menu

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.

LiveData engineeringReal data: UCI Online Retail II (CC BY 4.0)

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

  1. 1Extract

    Retries with backoff and jitter; structural errors caught per line

  2. 2Validate

    Data contract, business and integrity rules with reason codes, de-duplication

  3. 3Transform

    Normalise codes and text, classify lines, stable line keys

  4. 4Quality gate

    Batch checks; a failed gate stops the run before any write

  5. 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.

Three sources feed five stages: extract, validate, transform, quality gate and load. Rejected rows go to quarantine, published rows to the PostgreSQL warehouse, and every run records its metadata.

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.