Skip to content

Latest commit

 

History

2 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Modern Data Pipeline & Columnar OLAP Engine

CI Python 3.11+ DuckDB Polars FastAPI Docker Hugging Face License: MIT

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


🏗️ Architecture Overview

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
Loading

⚡ Key Engineering Features

  • 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 Dockerfile and docker-compose.yml with health-check monitoring and non-root execution.
  • Automated CI/CD: GitHub Actions test pipeline verifying unit tests across multiple Python runtimes.

🚀 Quickstart Guide

Option 1: Run with Docker Compose (Recommended)

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

The API and interactive documentation will be live at:

Option 2: Local Python Setup

# 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

📊 Analytical Benchmark

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

📡 API Endpoints Overview

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.

🧪 Testing & Validation

Run the comprehensive unit and integration test suite:

pytest tests/ -v --tb=short

The test suite validates:

  1. Synthetic stream generation statistical distributions.
  2. Vectorized Polars DataFrame schema casting.
  3. DuckDB SQL aggregations, quantile computations, and error boundary handling.
  4. FastAPI async HTTP endpoints, input contract validation, and response serialization.

👨‍💻 Author & Contact

Fabio Ignacio Torres Benítez
Data Engineer | Full-Stack & AI Systems Engineer

About

Modern Data Stack pipeline for high-throughput stream ingestion and sub-50ms columnar OLAP analytics using DuckDB, Polars, and FastAPI.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages