A pluggable Python ETL pipeline — reads CSV/Excel/JSON/Parquet/URL, validates the schema, transforms, enriches, and loads into Postgres/SQLite/Parquet with hash-based CDC, generates versioned snapshots, and ships with a Streamlit dashboard.
- 📥 Multi-format extract: CSV, TSV, Excel, JSON/JSONL, Parquet, HTTP (auto-detected)
- ✅ Schema validation with Pandera (fails fast if the input format drifts)
- 🧹 Transform with normalization (name, email, phone, value, date) and dedup by email/id
- 📊 Quality report: rejected rows saved with a reason to
rejected/ - 🔬 Enrichment: optional email MX validation + z-score anomaly detection
- 💾 Pluggable load: Postgres, SQLite, or Parquet
- 🔁 Modes:
replace,append,upsertwith hash-based CDC (only updates rows that actually changed) - 📈 Run metrics persisted to
etl_runs(rows in/out/rejected, duration, status) - 📸 Snapshots as timestamped Parquet files
- ⏱️ Scheduler via APScheduler (cron expression)
- 📊 Streamlit dashboard with charts for runs, data quality, data, and anomalies
- 🐳 Docker Compose with Postgres + ETL + dashboard
- 🧪 Tests with pytest
etl-local/
├── etl/
│ ├── cli.py # Typer CLI
│ ├── config.py # Pydantic settings (.env)
│ ├── pipeline.py # Orchestrator
│ ├── extract/ # CSV, Excel, JSON, Parquet, HTTP
│ ├── transform.py # Normalization + dedup
│ ├── validate.py # Pandera schemas
│ ├── enrich.py # MX + anomaly detection
│ ├── quality.py # Rejected-rows report
│ ├── load/ # Postgres, SQLite, Parquet
│ ├── metrics.py # etl_runs table
│ ├── snapshot.py # Parquet versions
│ ├── scheduler.py # APScheduler
│ └── logging_setup.py # Rich + rotating file handler
├── dashboard/app.py # Streamlit
├── tests/ # pytest
├── docker-compose.yml
├── Dockerfile
├── Makefile
├── pyproject.toml
└── .env.example
python -m venv venv
source venv/bin/activate
pip install -e .[dev]
cp .env.example .env # edit if you want PostgresMain CLI:
etl run # uses .env config
etl run --target sqlite --mode upsert # override via flags
etl run --input data.xlsx # another file
etl run --input https://example.com/data.csv
etl validate clientes.csv # validate schema only
etl runs # list recent runs
etl snapshots # list snapshots on disk
etl dashboard # launch Streamlit on :8501
etl schedule --cron "*/15 * * * *" # schedule periodic runsOr via Make:
make test # pytest
make run # pipeline
make dashboard # Streamlit
make schedule # foreground scheduler
make docker-up # Postgres + ETL + dashboard via docker compose# SQL target (when ETL_TARGET=postgres)
DB_HOST=localhost
DB_PORT=5433
DB_NAME=etldb
DB_USER=etl
DB_PASSWORD=etl123
# Pipeline
ETL_INPUT=clientes.csv
ETL_TARGET=sqlite # postgres | sqlite | parquet
ETL_TABLE=clientes
ETL_MODE=upsert # replace | append | upsert
ETL_SQLITE_PATH=etl.db
# Enrichment
ETL_VALIDATE_MX=false
ETL_DETECT_ANOMALIES=true
ETL_ANOMALY_ZSCORE=3.0
# Scheduler
ETL_CRON=0 */6 * * *- replace — drops and recreates the table. Use for full reloads.
- append — insert only; keeps full history.
- upsert — hashes each row; inserts new ones, updates rows that changed, and skips unchanged rows (CDC).
docker compose up --build
# ETL runs once, dashboard stays available at http://localhost:8501pytest # 13 tests covering extract/transform/validate/enrich/load/pipeline
pytest --cov=etl # with coverage
ruff check etl tests # lintMIT