A Python-based real-time data engineering pipeline that collects hourly air quality data from the Open-Meteo API for Nairobi and Mombasa, stores it in MongoDB, streams it via Apache Kafka, and loads it into DataStax Astra DB using the Data API for scalable storage and analysis.
# Real-Time Air Quality Data Pipeline (MongoDB → Kafka → Astra DB)
This project builds a **real-time data pipeline** that collects **air quality data** from the Open-Meteo Air Quality API for **Nairobi** and **Mombasa**, stores it in **MongoDB**, streams it through **Kafka**, and loads it into **DataStax Astra DB** using its **Data API**.
---
## Project Overview
### Components
- **Data Source**: Open-Meteo Air Quality API
- **Data Storage**: MongoDB Atlas
- **Data Streaming**: Apache Kafka
- **Data Warehouse**: DataStax Astra DB (NoSQL)
- **Language**: Python
---
## Project Structure
```
air_quality_pipeline/
│
├── kafka/ # (ignored in .gitignore)
├── producer.py # Fetches API data, stores in MongoDB, and sends to Kafka
├── consumer.py # Reads Kafka messages and loads into Astra DB
├── .env # Environment variables (excluded from Git)
├── .gitignore # Files/folders ignored by Git
└── README.md # Project documentation
````
---
## How It Works
1. **Producer Script (`producer.py`)**
- Fetches hourly air quality data for Nairobi and Mombasa.
- Inserts the latest record into **MongoDB Atlas**.
- Publishes the record to a **Kafka topic** named `air_quality`.
2. **Consumer Script (`consumer.py`)**
- Listens to the `air_quality` Kafka topic.
- Inserts incoming messages into **Astra DB** via the **Data API**.
---
## Prerequisites
Before running the project, ensure you have:
- Python 3.9 or later
- MongoDB Atlas cluster (free tier works fine)
- Apache Kafka running locally
- DataStax Astra DB account and a database with **Data API enabled**
- A `.env` file with credentials
---
## Environment Variables (`.env`)
Create a `.env` file in the project root with the following content:
```bash
# MongoDB Atlas connection
MONGO_URI="your_mongodb_connection_string"
# Kafka configuration
KAFKA_TOPIC=air_quality
KAFKA_SERVER=localhost:9092 …