Skip to content

Repository files navigation

FinFlow

Real-Time Financial Data Engineering Platform — explorable as a cinematic 3D world

CI Python FastAPI React Three.js Kafka API dbt Docker License: MIT

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.


Table of contents


Why this project exists

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 interactive 3D platform

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/Linux

Design notes and the backend↔frontend mapping: docs/3d_experience.md · decision record: ADR 0004.


Architecture

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]
Loading

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.


Tech stack

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

What it demonstrates

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

Quickstart

1. Full platform (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

2. Just the 3D app, locally (no Docker)

# 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/Linux

With 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).

3. Optional: LLM features

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.


Repository structure

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

Data model

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.


What breaks, and how it's handled

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_letters audit table. Never retried, never block the stream.
  • Transient failures (DB hiccup) — bounded exponential-backoff retries (tenacity); after MAX_DELIVERY_ATTEMPTS the 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.

Testing & CI

# 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 build

GitHub Actions runs on every push/PR: ruff, mypy, pytest, dbt parse, the frontend type-check + build, and Docker image/compose validation.


Roadmap

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

About

Real-time financial data engineering platform (streaming + medallion warehouse + quality gates), explorable as a cinematic 3D pipeline flythrough. FastAPI + React/Three.js over Kafka, dbt, Airflow, Spark, Postgres.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages