End-to-end data pipeline for pharmacy stock monitoring, demand forecasting, and real-time alerting.
Ingestion (Kafka) → Streaming (Spark) → Warehouse (BigQuery) → dBT → ML / RAG → Dashboard (Streamlit)
↕
Observability Layer
Stack: Kafka, Spark Structured Streaming, BigQuery, dBT, XGBoost, MLflow, ChromaDB, Streamlit, Airflow, Docker
├── ingestion/ Kafka producer + event simulator + Pydantic validation
├── streaming/ Spark consumer, anomaly detection, alert engine
├── warehouse/ BigQuery migrations + Terraform
├── dbt_models/ dBT transformations (staging → intermediate → marts)
├── ml/ XGBoost demand forecasting + MLflow registry
├── rag/ RAG chatbot (ChromaDB + sentence-transformers)
├── dashboard/ Streamlit dashboard (sales, inventory, forecasts, chatbot)
├── airflow/dags/ Airflow DAGs (dbt, ML retrain, vector update, data quality)
├── infra/docker/ Docker Compose + Dockerfiles for all services
├── tests/ Python tests (pytest)
└── observability/ Pipeline monitoring & alerting
- Docker & Docker Compose
- Python 3.10+
# Setup
python -m venv venv
venv\Scripts\activate # Windows
pip install -r requirements.txt
# Run simulator
python ingestion/simulator.py
# Run dashboard (reads mock data)
streamlit run dashboard/app.py
# ML training
python -m ml.demand_forecasting.train
# Tests
pytest tests/ -v# Full stack (Kafka, Spark, MLflow, Dashboard, Airflow)
docker-compose -f infra/docker/docker-compose.yml up -d
# Services:
# Dashboard: http://localhost:8501
# Airflow: http://localhost:8080 (admin/admin)
# MLflow: http://localhost:5000- Simulator generates realistic pharmacy events (sales, restocks, expiry, returns)
- Kafka Producer validates via Pydantic & publishes to 4 topics
- Spark Streaming consumes, detects anomalies (low stock, expiry alerts), writes to BigQuery
- dBT transforms raw data through staging → intermediate → marts
- ML trains XGBoost on mart data, logs to MLflow, generates 30-day forecasts
- RAG indexes drug data + forecasts into ChromaDB for chatbot
- Dashboard visualizes KPIs, forecasts, and chatbot
- Airflow orchestrates daily dbt runs, weekly ML retraining, vector index updates
- Fixed SQL syntax errors (trailing commas, wrong column refs) in dbt models
- Fixed Dockerfile paths and Docker Compose volume mounts
- Added missing
drugssource to dBTsources.yml - Fixed low_stock_job.py business logic (was inverted)
- Consolidated duplicate Dockerfiles to
infra/docker/ - Standardized Airflow version to 2.10.x
- Removed dead dependencies (langchain, pinecone)
- Added 21 Python tests (ingestion schemas, ML metrics, feature engineering)
- Fixed 70+ lint issues via ruff
- Created checkpoints/ directory for Spark streaming
- Added experiment None-check in ML retrain DAG
See .env.example for required variables:
GCP_PROJECT_ID— BigQuery projectGOOGLE_APPLICATION_CREDENTIALS— Service account pathKAFKA_BOOTSTRAP_SERVERS— Kafka brokerSLACK_WEBHOOK_URL— Alert channel
# Run all tests
pytest tests/ -v
# Lint
ruff check .
# dBT tests
cd dbt_models && dbt test