ETL β PostgreSQL β dbt β Airflow orchestration. Built on real OpenWeather + GeoNames data across 209 countries, 583 cities, 88 features.
| Layer | Tool |
|---|---|
| Extract & Transform | Python 3.11, Pandas |
| Load | SQLAlchemy + PostgreSQL (raw schema) |
| Transformation | dbt (staging β marts) |
| Orchestration | Apache Airflow 2.x |
| Visualisation | Power BI Desktop (via ODBC) |
| Environment | Linux (Ubuntu) backend, Windows analytics layer |
| Next | Docker, GitHub Actions |
gaqip/
βββ main.py # Pipeline entry point
βββ extract/ # API + CSV ingestion
βββ transform/ # Pandas transformations
βββ load/ # SQLAlchemy DB loader
βββ air_quality_project/ # dbt project root
β βββ models/
β β βββ staging/ # stg_weather_air_quality
β β βββ marts/ # mart_city_aqi, mart_country_summary
β βββ dbt_project.yml
β βββ profiles.yml
βββ dags/
β βββ gaqip_dag.py # Airflow DAG
βββ logs/
β βββ pipeline.log
βββ requirements.txt
python main.pyRuns full Extract β Transform β Load into raw.weather_air_quality.
cd air_quality_project
dbt run
dbt testBuilds staging and mart models. Tests enforce schema and data quality constraints.
DAG: gaqip_dag.py
gaqip_pipeline
βββ run_etl (BashOperator β sources venv, runs main.py)
βββ dbt_run (BashOperator β dbt run)
βββ dbt_test (BashOperator β dbt test)
βββ notify (success/failure alert)
default_args = {
"retries": 2,
"retry_delay": timedelta(minutes=5),
}Dataset is a non-incremental snapshot β no append accumulation.
# Pattern: TRUNCATE β LOAD
with engine.begin() as conn:
conn.execute(text("TRUNCATE TABLE raw.weather_air_quality"))
df.to_sql(
name="weather_air_quality",
con=conn,
schema="raw",
if_exists="append",
index=False
)Guarantees deterministic state on every run. Safe for dbt model rebuilds.
# 1. Clone and create venv
git clone https://github.com/<you>/gaqip.git && cd gaqip
python3.11 -m venv gp_venv && source gp_venv/bin/activate
pip install -r requirements.txt
# 2. Configure DB
export DB_URL="postgresql://user:pass@localhost:5432/gaqip"
# 3. Run pipeline
python main.py
# 4. Run dbt
cd air_quality_project && dbt run && dbt test
# 5. (Optional) Start Airflow
airflow db init
airflow scheduler & airflow webserverAirflow tasks use
source gp_venv/bin/activateinside BashOperators β not the system Python.
| Layer | Strategy |
|---|---|
| ETL | try/except with retry decorator, logged to logs/pipeline.log |
| DB load | engine.begin() β full transactional commit, auto-rollback on error |
| Airflow | Task-level retry (retries: 2, retry_delay: 5m) |
| dbt | Model-level dbt test β fails DAG on assertion breach |
1. ResourceClosedError on DB load
- Root cause: connection used outside transaction scope
- Fix: enforce
engine.begin()context manager throughout loader
2. Airflow task using wrong Python environment
- Root cause: BashOperator defaulting to system Python
- Fix: explicit
source gp_venv/bin/activatein every BashOperator command
3. dbt schema drift (staging_mart vs staging_marts)
- Root cause: conflicting schema names between
dbt_project.ymland profile - Fix: centralised schema config in
dbt_project.yml, removed profile-level overrides
4. Power BI native connector SSL failure
- Root cause: Power BI's .NET-based PostgreSQL driver enforces strict TLS certificate validation; the PostgreSQL server used a self-signed certificate, causing rejection at the TLS handshake β not a credentials or network issue
- pgAdmin and
psqlboth connected successfully (both accept self-signed certs by default) - Fix: bypassed native connector entirely via PostgreSQL ODBC driver (
psqlodbc_x64.msi)
# ODBC DSN configuration (Windows ODBC Data Source Administrator β 64-bit)
Driver: PostgreSQL Unicode(x64)
Host: <linux-server-ip>
Port: 5432
Database: gaqip
SSL Mode: disable β bypasses TLS handshake enforcement
Power BI connection path:
Get Data β ODBC β Select DSN β Load tables
| Tool | SSL behaviour |
|---|---|
| pgAdmin | Accepts self-signed certificates |
| psql CLI | Accepts self-signed certificates |
| Power BI native connector | Strict TLS validation β rejects self-signed β |
| Power BI via ODBC | Configurable SSL mode β connection stable β |
Logs written per stage to logs/pipeline.log:
[EXTRACT] 583 rows fetched β 2026-04-24T07:06:06Z
[TRANSFORM] Nulls filled: wind_gust=47, grnd_level=12
[LOAD] TRUNCATE β INSERT 583 rows β OK
[DBT] dbt run completed: 4 models, 0 errors
[DBT] dbt test completed: 11 tests passed
staging/
βββ stg_weather_air_quality β typed, null-handled, renamed cols
marts/
βββ mart_city_aqi β city-level AQI ranking + pm2_5 category
βββ mart_country_summary β country-level aggregations (avg temp, AQI dist)
| Column | Type | Description |
|---|---|---|
country_code |
str | ISO 2-letter code |
city_name |
str | GeoNames city |
temp |
float | Β°C (OpenWeather metric) |
aqi |
int | 1β5 (OpenWeather AQI scale) |
pm2_5, pm10 |
float | Β΅g/mΒ³ |
target_aqi_class |
int | Encoded AQI label (ML target) |
comfort_proxy |
float | Derived: temp Γ (1 - humidity/100) |
wind_pm25_interaction |
float | wind_speed Γ pm2_5 |
Full schema: 88 columns, 583 rows, 209 countries, 5 continents.
| Item | Status |
|---|---|
| Modular ETL | β |
| Transactional DB load | β |
| Retry logic (ETL + Airflow) | β |
| Structured logging | β |
| dbt staging + marts | β |
| dbt tests | β |
| Linux execution stability | β |
| Airflow DAG (BashOperator) | β |
| Power BI integration (ODBC) | β |
| Dockerization | β³ next |
| CI/CD (GitHub Actions) | β³ next |
| PythonOperator DAG refactor | β³ next |
| FastAPI data layer | β³ planned |
This project was built through real production failures β not a tutorial:
- Migrated from Windows β Linux for Airflow stability
- Resolved Python 3.12 incompatibility with Airflow (pinned to 3.11)
- Debugged
ResourceClosedErrorfrom connection lifecycle mismanagement - Fixed dbt schema drift from profile/project config conflict
- Stabilised Airflow venv isolation in BashOperator tasks
- Resolved Power BI SSL rejection via ODBC driver (psqlodbc_x64) with
sslmode=disable
- Dockerize β Postgres + Airflow + dbt in
docker-compose.yml - CI/CD β GitHub Actions: lint β test β dbt run on push
- PythonOperator refactor β replace BashOperators with native Python callables
- Great Expectations β row-level data quality assertions pre-load
- FastAPI layer β
/cities,/aqi,/countriesendpoints over the mart layer