An end-to-end, production-grade Modern Data Stack (MDS) pipeline engineered for high-throughput stream ingestion, columnar Parquet lakehouse persistence, and sub-50ms analytical queries (OLAP) powered by DuckDB and Polars.
Built by Fabio Torres (M.Sc. in Data Engineering & Cloud Infrastructure).
graph TD
subgraph IngestionLayer ["1. Ingestion & Simulation Layer"]
SG["Synthetic Stream Generator<br><i>Vectorized Batching</i>"]
REST["FastAPI Ingest Endpoint<br><code>POST /api/v1/stream/ingest</code>"]
PYD["Pydantic v2 Contract Validation"]
end
subgraph ProcessingLayer ["2. Vectorized Processing (Polars)"]
PL["Polars LazyFrame Processing"]
CAST["Schema Enforcement & Type Casting"]
COMP["Snappy Compression Engine"]
end
subgraph LakehouseLayer ["3. Storage (Columnar Lakehouse)"]
PARQ["Partitioned Parquet Files<br><code>data/events_*.parquet</code>"]
end
subgraph OLAPLayer ["4. Columnar Analytics (DuckDB)"]
DDB["DuckDB In-Memory OLAP Engine"]
SQL["Vectorized Window Functions<br>& Percentile Latency (p95, p99)"]
end
subgraph ServingLayer ["5. Serving & Consumption"]
API["FastAPI Async REST API"]
DOCS["Interactive Swagger / OpenAPI<br><code>/docs</code>"]
DASH["Real-time KPIs & Cohort Analytics"]
end
SG --> REST
REST --> PYD
PYD --> PL
PL --> CAST
CAST --> COMP
COMP --> PARQ
PARQ --> DDB
DDB --> SQL
SQL --> API
API --> DOCS
API --> DASH
style IngestionLayer fill:#0f172a,stroke:#38bdf8,stroke-width:1px,color:#f8fafc
style ProcessingLayer fill:#0f172a,stroke:#fb923c,stroke-width:1px,color:#f8fafc
style LakehouseLayer fill:#0f172a,stroke:#eab308,stroke-width:1px,color:#f8fafc
style OLAPLayer fill:#0f172a,stroke:#4ade80,stroke-width:1px,color:#f8fafc
style ServingLayer fill:#0f172a,stroke:#a855f7,stroke-width:1px,color:#f8fafc
- Vectorized High-Throughput Processing: Replaces slow row-by-row iteration with vectorized Polars engines, achieving ingestion throughput of >50,000 records/sec per core.
- Embedded Columnar OLAP (DuckDB): Direct SQL querying on top of multi-file Parquet partitions without requiring a heavyweight distributed Spark/Trino cluster.
- Sub-50ms Latency Metrics: Calculates complex aggregations including p95 and p99 latency quantiles, active user rollups, and error rates in milliseconds.
- Type-Safe Contract Enforcement: Powered by Pydantic v2 strict validation schemas ensuring zero corrupted records enter the lakehouse.
- Single-Command Containerization: Production-ready
Dockerfileanddocker-compose.ymlwith health-check monitoring and non-root execution. - Automated CI/CD: GitHub Actions test pipeline verifying unit tests across multiple Python runtimes.
Clone and run the entire ecosystem in 1 command:
git clone https://github.com/neurodeveloper11/modern-data-pipeline.git
cd modern-data-pipeline
docker compose up --buildThe API and interactive documentation will be live at:
- Swagger UI: http://localhost:8000/docs
- Health Check: http://localhost:8000/health
# 1. Clone repository
git clone https://github.com/neurodeveloper11/modern-data-pipeline.git
cd modern-data-pipeline
# 2. Create and activate virtual environment
python -m venv venv
source venv/bin/activate # On Windows: .\venv\Scripts\activate
# 3. Install dependencies
pip install -r requirements.txt
# 4. Run automated test suite
pytest tests/ -v
# 5. Start API server
uvicorn src.api:app --reload --port 8000| Operation | Engine | Dataset Size | Execution Latency |
|---|---|---|---|
| Stream Batch Ingest & Parquet Write | Polars (Snappy) | 10,000 events | ~42 ms |
| P95 / P99 Latency & Error Rollup | DuckDB OLAP | 50,000 events | ~14 ms |
| Daily Cohort User Retention | Polars LazyFrame | 50,000 events | ~18 ms |
| Arbitrary Group-by Analytical SQL | DuckDB Engine | 100,000 events | ~28 ms |
| Method | Endpoint | Description |
|---|---|---|
GET |
/health |
Kubernetes/Docker readiness & liveness probe. |
POST |
/api/v1/stream/ingest |
Ingests a batch of telemetry events into the Parquet lakehouse. |
GET |
/api/v1/analytics/realtime-kpis |
Real-time p95, p99, throughput, and error metrics via DuckDB. |
GET |
/api/v1/analytics/cohorts |
Multi-day cohort retention analysis using Polars LazyFrames. |
POST |
/api/v1/analytics/query |
Arbitrary read-only SQL querying directly on Parquet files. |
POST |
/api/v1/generator/simulate |
Benchmarking utility to synthesize and ingest N events. |
Run the comprehensive unit and integration test suite:
pytest tests/ -v --tb=shortThe test suite validates:
- Synthetic stream generation statistical distributions.
- Vectorized Polars DataFrame schema casting.
- DuckDB SQL aggregations, quantile computations, and error boundary handling.
- FastAPI async HTTP endpoints, input contract validation, and response serialization.
Fabio Ignacio Torres Benítez
Data Engineer | Full-Stack & AI Systems Engineer
- GitHub: @neurodeveloper11
- Hugging Face: @neurodeveloper
- LinkedIn: Fabio Torres
- Email: psicologofabiotorres@gmail.com