An end-to-end platform that ingests real-time crypto trades (Binance WebSocket) and daily equity prices (yfinance), processes them through a streaming + medallion-warehouse architecture with an enforced data-quality gate, and serves the result two ways: a classic trading terminal and an interactive 3D flythrough where the world is the architecture.
Built entirely on free tooling. Runs with a single docker compose up — or with zero
Docker via a bundled demo database.
- Why this project exists
- The interactive 3D platform
- Architecture
- Tech stack
- What it demonstrates
- Quickstart
- Repository structure
- Data model
- What breaks, and how it's handled
- Testing & CI
- Roadmap
A data-engineering portfolio lives or dies on whether a reviewer can grasp the architecture in the first minute. Most projects bury it in a README diagram that's easy to skim past. FinFlow makes the pipeline the thing you experience: a production-grade platform you can fly through, inspecting each stage's live metrics, the technology it uses, and how it handles failure.
The stack deliberately maps to the most-requested Data Engineer skills — Python + SQL, Kafka-style streaming, Airflow orchestration, dbt warehouse modeling, PySpark, Docker, CI/CD, cloud + IaC, and GDPR-aware governance — plus a genuine full-stack slice (typed API, WebSockets, a TypeScript/Three.js SPA) on top.
Hard constraints honoured: everything is free (public crypto WebSockets, yfinance, Groq free tier, GCP BigQuery sandbox), everything runs via Docker Compose, and the data-quality gate runs before load — a failed checkpoint stops the pipeline.
The headline experience. A cinematic, navigable flythrough of the pipeline itself, built as a full-stack app: FastAPI backend + React / TypeScript / Three.js frontend.
Live market data streams as glowing packets from ingestion, through the Kafka-style bus and the quality gate, up the bronze → silver → gold warehouse (the spine literally rises as data is refined), and out to analytics. It opens with an auto-playing guided tour, then hands you the controls to orbit freely and click into any stage.
- Every node is a real stage — with live numbers, the tech it uses, the data-engineering skill it demonstrates, and its failure-handling story (DLQ, idempotency, quality gate).
- Live candlesticks and an anomaly feed are embedded directly in the node panels, served from the gold marts over a real REST + WebSocket API.
- Runs with or without Docker. The API auto-detects the Postgres warehouse and falls back to a bundled SQLite demo database with a synthetic live-tick feed — so the demo is always live.
# Full platform (live warehouse) # Or instantly, zero Docker (demo feed)
docker compose up -d --build ./scripts/serve.ps1 # Windows → :8000
# → http://localhost:8090 ./scripts/serve.sh # macOS/LinuxDesign notes and the backend↔frontend mapping: docs/3d_experience.md · decision record: ADR 0004.
flowchart LR
subgraph Ingestion
WS[Binance WebSocket<br/>real-time trades] --> P[Producer]
YF[yfinance<br/>daily OHLCV] --> AF[Airflow DAG]
end
P --> RP[(Redpanda<br/>market.trades.raw)]
RP --> C[Consumer<br/>idempotent upserts]
C -- "poison / exhausted retries" --> DLQ[(Dead-letter topic<br/>+ audit table)]
C --> BR[(Postgres · bronze)]
AF --> GE{Great Expectations<br/>quality gate}
GE -- pass --> BR
GE -- fail --> STOP[pipeline stops]
BR --> DBT[dbt<br/>bronze → silver → gold]
BR --> SP[PySpark<br/>daily aggregates]
SP --> GOLD[(Postgres · gold marts)]
DBT --> GOLD
GOLD --> API[FastAPI<br/>REST + WebSocket]
API --> WEB[React + Three.js<br/>3D platform]
GOLD --> ST[Streamlit<br/>trading terminal]
GOLD --> LLM[Groq LLM<br/>anomaly explain + NL→SQL]
LLM --> API
LLM --> ST
P -. metrics .-> PROM[Prometheus]
C -. metrics .-> PROM --> GRAF[Grafana]
The serving layer presents the same gold marts two ways: the Streamlit terminal and the full-stack 3D platform. Full design details in docs/architecture.md.
| Layer | Tool | Notes |
|---|---|---|
| Language | Python 3.11+ · TypeScript | Type hints / strict TS throughout |
| Real-time ingestion | Binance public WebSocket | Free, no key, genuine real-time ticks |
| Batch ingestion | yfinance | Free, no key; daily equity OHLCV |
| Streaming bus | Redpanda (Kafka API) | Same skills as Kafka, less RAM |
| Stream processing | confluent-kafka consumer | Idempotent upserts + dead-letter topic |
| Batch processing | PySpark | Daily VWAP / aggregates via JDBC |
| Orchestration | Apache Airflow | Batch DAGs, retries, backfills |
| Transformation | dbt Core (dbt-postgres) |
Medallion layering, tests, lineage |
| Warehouse | PostgreSQL | Optional BigQuery sandbox branch |
| Data quality | Great Expectations + dbt tests | Hard gate before load |
| Observability | Prometheus + Grafana | Lag, throughput, freshness, DLQ depth |
| Backend API | FastAPI + WebSockets | Serves topology + live data to the 3D app |
| Frontend | React · Three.js (react-three-fiber) | Cinematic 3D pipeline flythrough |
| Analytics UI | Streamlit (lightweight-charts) | Classic trading terminal |
| LLM (optional) | Groq free tier | Anomaly explain + guarded NL→SQL |
| IaC | Docker Compose + Terraform (GCP) | BigQuery dataset + GCS bucket |
| CI/CD | GitHub Actions | ruff · mypy · pytest · dbt parse · web build · docker |
| Skill | Where |
|---|---|
| Real-time streaming (Kafka API) | ingestion/streaming/ — WebSocket → Redpanda → Postgres |
| Failure handling | Idempotent upserts, bounded retries, dead-letter queue + audit table |
| Orchestration | airflow/dags/finflow_batch.py — extract → quality gate → load → dbt |
| Data quality as a hard gate | quality/ — Great Expectations before load; failure stops the run |
| Analytics engineering | dbt/ — medallion layering, schema tests, lineage docs |
| Distributed batch processing | spark/jobs/daily_agg.py — PySpark VWAP/aggregates via JDBC |
| Observability | observability/ — Prometheus metrics, Grafana dashboards |
| Full-stack engineering | api/ (FastAPI + WS) and web/ (React + Three.js) over the warehouse |
| LLM integration, safely | llm/ — NL→SQL behind a unit-tested read-only allowlist guard |
| GDPR awareness | governance/ — PII detection, deterministic pseudonymization, log scrubbing |
| IaC + cloud | terraform/ — BigQuery (sandbox) + GCS, region europe-west3 |
| CI/CD | .github/workflows/ci.yml — lint, types, tests, dbt parse, web build, docker |
cp .env.example .env # defaults work out of the box
docker compose up -d --build| Service | URL |
|---|---|
| FinFlow 3D Platform (React + Three.js) | http://localhost:8090 |
| FinFlow Terminal (Streamlit) | http://localhost:8501 |
| Airflow | http://localhost:8080 (admin/admin) |
| Grafana ops dashboards | http://localhost:3000 (admin/admin) |
| Redpanda Console (topics, DLQ) | http://localhost:8085 |
| Prometheus | http://localhost:9090 |
Within ~30 s the producer is streaming live trades and the consumer is writing them to
bronze.trades. Then run the batch + transformation layers:
# batch pipeline: equities → quality gate → load → dbt
docker compose exec airflow airflow dags trigger finflow_batch
# or run dbt directly
docker compose exec airflow bash -c "cd /opt/finflow/dbt && dbt build --profiles-dir ."
# submit the Spark job
docker compose exec spark spark-submit \
--packages org.postgresql:postgresql:42.7.4 /opt/spark-jobs/daily_agg.py# Hot-reload dev (FastAPI :8000 + Vite :5173)
./scripts/dev.ps1 # Windows → http://localhost:5173
./scripts/dev.sh # macOS/Linux
# Single-origin serve (build once; FastAPI serves the SPA + API)
./scripts/serve.ps1 # Windows → http://localhost:8000
./scripts/serve.sh # macOS/LinuxWith no warehouse running, the API serves the bundled SQLite demo DB and a synthetic live tick
feed — fully interactive offline. Start docker compose up -d and it switches to the live
Postgres warehouse automatically (FINFLOW_DATA_SOURCE=auto).
Anomaly explanations and natural-language → SQL use Groq's free tier. Grab a key at
console.groq.com and set GROQ_API_KEY in .env.
finflow/
├── api/ # FastAPI backend: topology + live data + WS feed (serves the SPA)
├── web/ # React + TypeScript + Three.js 3D frontend
│ ├── src/scene/ # 3D scene: nodes, flows, camera rig, environment
│ └── src/hud/ # DOM overlay: intro, tour, node panels, live charts
├── ingestion/
│ ├── streaming/ # WebSocket producer + idempotent Kafka consumer (DLQ)
│ └── batch/ # yfinance extractors
├── airflow/dags/ # batch orchestration DAGs
├── spark/ # PySpark batch jobs
├── dbt/ # models: bronze / silver / gold + schema tests
├── quality/ # Great Expectations suites + checkpoints (the gate)
├── llm/ # Groq anomaly explanation + guarded NL→SQL
├── dashboard/ # Streamlit trading terminal
├── observability/ # Prometheus config + Grafana dashboards
├── governance/ # PII detection / pseudonymization
├── terraform/ # optional GCP (BigQuery + GCS) module
├── demo/ # SQLite demo DB + seed script (zero-Docker fallback)
├── scripts/ # dev.ps1/.sh, serve.ps1/.sh — local run helpers
├── tests/ # pytest unit + integration
├── docs/ # architecture, ADRs, data dictionary, 3D design
├── docker/ # per-service Dockerfiles (incl. docker/app for the 3D platform)
├── docker-compose.yml
└── .github/workflows/ # CI
Medallion architecture in PostgreSQL:
| Layer | Schema | Contents | Materialization |
|---|---|---|---|
| Bronze | bronze |
Raw, 1:1 with source; idempotent natural-key upserts | Physical tables |
| Silver | silver |
Cleaned, typed, deduplicated | dbt views |
| Gold | gold |
Business marts: 1-min OHLCV, daily equity + SMAs, z-score anomalies, Spark aggregates | dbt tables / incremental |
Full column-level reference: docs/data_dictionary.md.
Production pipelines are defined by their failure modes, not their happy paths:
- Poison messages (malformed JSON, wrong event type) — parsed with a strict schema; failures
route immediately to the dead-letter topic and a queryable
bronze.dead_lettersaudit table. Never retried, never block the stream. - Transient failures (DB hiccup) — bounded exponential-backoff retries (
tenacity); afterMAX_DELIVERY_ATTEMPTSthe message goes to the DLQ and the stream continues. - Duplicate delivery (Kafka is at-least-once) — every write is an idempotent upsert on a natural key; replays are no-ops by design, verified by tests.
- Bad source data — the Great Expectations gate runs before load; a failed checkpoint stops the DAG. Bad batches never reach the warehouse.
- WebSocket drops — auto-reconnect with exponential backoff; reconnects are a Prometheus counter. The frontend tick socket reconnects with bounded backoff too.
- LLM misbehaviour — generated SQL passes a unit-tested guard (single statement, SELECT-only, table allowlist) before it touches the warehouse.
- Warehouse unavailable — the API degrades gracefully to the bundled demo DB so the presentation layer never goes dark.
# Python
pip install -r requirements.txt -r requirements-dev.txt
ruff check . && mypy ingestion quality llm governance && pytest tests/
# Frontend
cd web && npm install && npm run build # type-check + production buildGitHub Actions runs on every push/PR: ruff, mypy, pytest, dbt parse, the frontend type-check + build, and Docker image/compose validation.
Built bottom-up, one layer at a time (infra → streaming slice → resilience → batch + quality → dbt → Spark → observability → LLM → governance/cloud → the full-stack 3D platform). Next:
- Capture screenshots / a short clip of the 3D platform and the Grafana boards
- Wire the live WebSocket tick stream into per-frame packet bursts in the 3D scene
- BigQuery sandbox branch via the Terraform module
- Code-split the Three.js bundle
Built by SaadH-077 · MIT licensed · Free tooling only