This project demonstrates a production-style batch data pipeline built using a Bronze → Silver → Gold architecture. It ingests real-time weather data, processes it into structured formats, and loads it into a PostgreSQL data warehouse for analytics.
The pipeline is designed to mimic real-world data engineering workflows, including incremental loading, partitioned data lakes, and SQL-based analytics readiness.
Weather API
↓
Python Ingestion Script
↓
S3 Bronze (Raw JSON - Immutable)
↓
Python Transformation (Pandas + PyArrow)
↓
S3 Silver (Partitioned Parquet)
↓
Python Incremental Loader
↓
PostgreSQL Gold (Analytics Layer)
- Programming Language: Python
- Libraries: Pandas, PyArrow, Requests, Boto3, SQLAlchemy, psycopg2, s3fs
- Storage: AWS S3 (Data Lake)
- Database: PostgreSQL (Local Data Warehouse)
- Concepts: ETL/ELT, Data Lake, Partitioning, Incremental Loading, SQL Analytics
-
Weather data fetched from a public weather API (Open-Meteo)
-
Data includes:
- Temperature
- Wind speed & direction
- Weather codes
- Observation timestamps
-
API data fetched using Python (
requests) -
Metadata enrichment added:
ingestion_timestamprun_id(UUID)source_system
-
Stored as immutable JSON files in S3
Structure:
s3://weather-data-raw-bhanu/
raw/weather/{run_id}.json
-
Raw JSON processed using Pandas
-
Extracted structured schema:
run_id ingestion_time temperature windspeed winddirection weathercode observation_time -
Converted to Parquet format using PyArrow
-
Stored in S3 with event-time partitioning
Partitioned Structure:
s3://weather-data-processed-bhanuu/silver/weather/
year=YYYY/
month=MM/
day=DD/
weather_HHMMSS.parquet
- PostgreSQL used as analytics layer
- Data loaded directly from S3 (no local storage)
| Column | Type |
|---|---|
| observation_time | TIMESTAMP (PK) |
| temperature | DOUBLE |
| windspeed | DOUBLE |
| winddirection | INTEGER |
| weathercode | INTEGER |
| ingestion_time | TIMESTAMP |
| run_id | TEXT |
Implemented watermark-based incremental loading:
-
Query latest record from PostgreSQL:
SELECT MAX(observation_time) FROM weather_observations;
-
Filter new records in Python:
df = df[df["observation_time"] > last_time]
-
Insert only new data
- Prevents duplicate data
- Ensures idempotent pipeline runs
- Improves performance
project-root/
│
├── ingestion/
│ └── ingestion.py
│
├── transformation/
│ └── transformation.py
│
├── gold/
│ └── load_to_postgre.py
│
├── logs/
│
└── README.md
pip install pandas pyarrow boto3 s3fs sqlalchemy psycopg2-binaryaws configurepython ingestion/ingestion.py
python transformation/transformation.py
python gold/load_to_postgre.py- Built a complete end-to-end data pipeline
- Used data lake + data warehouse architecture
- Implemented event-time partitioning
- Enabled direct S3 → PostgreSQL loading
- Designed incremental batch processing system
- Ensured idempotent and production-safe execution
- Add orchestration (Airflow / cron jobs)
- Build analytical tables (daily summaries, trends)
- Add data quality checks
- Dockerize the pipeline
- Integrate dashboarding (Power BI / Tableau)
This project demonstrates:
- Real-world data engineering workflows
- Handling semi-structured → structured data
- Working with cloud storage (S3)
- Designing scalable data pipelines
- Implementing incremental loading strategies
- Preparing data for analytics
Bhanu Prasad Aspiring Data Engineer
Give it a star ⭐ and feel free to fork or contribute!