- name
- realtime-cinema-data-engineering-pipeline
- description
- Build end-to-end real-time data pipelines with Kafka, PostgreSQL, Airflow, and Streamlit using Medallion Architecture for streaming analytics.
- triggers
- ["set up a real-time data pipeline with kafka and airflow","build a streaming etl pipeline with medallion architecture","create a cinema analytics dashboard with streamlit","implement bronze silver gold data layers","stream events through kafka to postgresql","orchestrate elt pipelines with apache airflow","visualize real-time analytics with plotly","configure kafka producers and consumers for data ingestion"]
# CinéWorld Real-Time Data Engineering Pipeline Skill
> Skill by [ara.so](https://ara.so) — Data Skills collection.
## Overview
This project implements an end-to-end real-time data engineering pipeline using Apache Kafka for event streaming, PostgreSQL for data warehousing with Medallion Architecture (Bronze/Silver/Gold layers), Apache Airflow for ELT orchestration, and Streamlit for live visualization. Perfect for learning how to build production-grade streaming data pipelines that process 1M+ events.
## Installation
### Prerequisites
- Docker and Docker Compose
- Python 3.8+
- Virtual environment (recommended)
### Setup Steps
```bash
# Clone the repository
git clone https://github.com/BaidaneAyoub/realtime-cinema-data-engineering.git
cd realtime-cinema-data-engineering
# Create and activate virtual environment
python -m venv myenv
source myenv/bin/activate # On Windows: myenv\Scripts\activate
# Install dependencies
pip install -r requirements.txt
# Start infrastructure (Kafka, PostgreSQL, Airflow)
docker-compose up -d
```
**Important**: Wait 2-3 minutes for Airflow to fully initialize before proceeding.
## Architecture Components
### 1. Medallion Architecture Layers
**Bronze Layer**: Raw JSON event data ingested from Kafka
```sql
-- Bronze table stores raw events
CREATE TABLE bronze_transactions (
id SERIAL PRIMARY KEY,
raw_data JSONB NOT NULL,
ingested_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
```
**Silver Layer**: Normalized 3NF tables (Customers, Movies, Showtimes, Transactions)
```sql
-- Normalized dimension and fact tables
CREATE TABLE silver_customers (...);
CREATE TABLE silver_movies (...);
CREATE TABLE silver_showtimes (...);
CREATE TABLE silver_transactions (...);
```
**Gold Layer**: Materialized views for analytics
```sql
-- Business-ready aggregated data
CREATE MATERIALIZED VIEW gold_cinema_analytics AS
SELECT ...
```
### 2. Kafka Event Producer
Generate and stream synthetic cinema transaction events:
```python
# producer/main_producer.py
from kafka import KafkaProducer
from faker import Faker
import json
import time
import os
fake = Faker()
# Initialize Kafka producer
producer = KafkaProducer(
bootstrap_servers=os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092'),
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
def generate_ticket_sale():
"""Generate a synthetic ticket sale event"""
return {
"transaction_id": fake.uuid4(),
"customer": {
"customer_id": fake.uuid4(),
"name": fake.name(),
"email": fake.email(),
"phone": fake.phone_number()
},
"movie": {
"movie_id": fake.uuid4(),
"title": fake.catch_phrase(),
"genre": fake.random_element(['Action', 'Comedy', 'Drama', 'Horror']),
"duration_minutes": fake.random_int(90, 180)
},
"showtime": {
"showtime_id": fake.uuid4(),
"cinema_location": fake.city(),
"screen_number": fake.random_int(1, 10),
"showtime": fake.date_time_this_month().isoformat()
},
"payment": {
"amount": round(fake.random.uniform(8.0, 25.0), 2),
"payment_method": fake.random_element(['Credit Card', 'Cash', 'Gift Card']),
"currency": "USD"
},
"seats": [f"{fake.random_element(['A','B','C','D'])}{fake.random_int(1,20)}"],
"timestamp": fake.date_time_this_month().isoformat()
}
# Stream events continuously
def start_streaming(topic='cinema_transactions', interval=0.5):
"""Start producing events to Kafka topic"""
print(f"🎬 Starting producer on topic: {topic}")
try:
while True:
event = generate_ticket_sale()
producer.send(topic, value=event)
print(f"✅ Sent: {event['transaction_id']}")
time.sleep(interval)
except KeyboardInterrupt:
print("\n🛑 Stopping producer...")
finally:
producer.flush()
producer.close()
if __name__ == "__main__":
start_streaming()
```
### 3. Kafka Consumer (Bronze Layer Ingestion)
Consume events from Kafka and insert into PostgreSQL Bronze layer:
```python
# consumer/main_consumer.py
from kafka import KafkaConsumer
import psycopg2
from psycopg2.extras import Json
import json
import os
# Kafka consumer configuration
consumer = KafkaConsumer(
'cinema_transactions',
bootstrap_servers=os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092'),
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id='cinema-consumer-group',
value_deserializer=lambda x: json.loads(x.decode('utf-8'))
)
# PostgreSQL connection
def get_db_connection():
return psycopg2.connect(
host=os.getenv('POSTGRES_HOST', 'localhost'),
port=os.getenv('POSTGRES_PORT', '5432'),
database=os.getenv('POSTGRES_DB', 'cinema_dw'),
user=os.getenv('POSTGRES_USER', 'postgres'),
password=os.getenv('POSTGRES_PASSWORD', 'postgres')
)
def insert_bronze(connection, event_data):
"""Insert raw event into bronze layer"""
with connection.cursor() as cursor:
cursor.execute(
"""
INSERT INTO bronze_transactions (raw_data)
VALUES (%s)
""",
(Json(event_data),)
)
connection.commit()
# Start consuming
def start_consuming():
print("🎧 Starting consumer...")
conn = get_db_connection()
try:
for message in consumer:
event = message.value
insert_bronze(conn, event)
print(f"✅ Inserted: {event.get('transaction_id', 'unknown')}")
except KeyboardInterrupt:
print("\n🛑 Stopping consumer...")
finally:
conn.close()
consumer.close()
if __name__ == "__main__":
start_consuming()
```
### 4. Airflow ELT DAG (Bronze → Silver → Gold)
Orchestrate data transformation pipeline:
```python
# dags/bronze_to_silver_and_gold_elt.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from datetime import datetime, timedelta
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2024, 1, 1),
'email_on_failure': False,
'retries': 1,
'retry_delay': timedelta(minutes=5),
}
def bronze_to_silver():
"""Extract from Bronze, transform, and load into Silver layer"""
pg_hook = PostgresHook(postgres_conn_id='postgres_cinema_dw')
conn = pg_hook.get_conn()
cursor = conn.cursor()
# Fetch unprocessed bronze records (limit 50k per run)
cursor.execute("""
SELECT id, raw_data
FROM bronze_transactions
WHERE processed = FALSE
LIMIT 50000
""")
records = cursor.fetchall()
for record_id, raw_data in records:
# Extract nested fields
customer = raw_data['customer']
movie = raw_data['movie']
showtime = raw_data['showtime']
payment = raw_data['payment']
# Insert into silver_customers (upsert)
cursor.execute("""
INSERT INTO silver_customers (customer_id, name, email, phone)
VALUES (%(customer_id)s, %(name)s, %(email)s, %(phone)s)
ON CONFLICT (customer_id) DO NOTHING
""", customer)
# Insert into silver_movies (upsert)
cursor.execute("""
INSERT INTO silver_movies (movie_id, title, genre, duration_minutes)
VALUES (%(movie_id)s, %(title)s, %(genre)s, %(duration_minutes)s)
ON CONFLICT (movie_id) DO NOTHING
""", movie)
# Insert into silver_showtimes (upsert)
cursor.execute("""
INSERT INTO silver_showtimes (showtime_id, cinema_location, screen_number, showtime)
VALUES (%(showtime_id)s, %(cinema_location)s, %(screen_number)s, %(showtime)s)
ON CONFLICT (showtime_id) DO NOTHING
""", showtime)
# Insert into silver_transactions (fact table)
cursor.execute("""
INSERT INTO silver_transactions
(transaction_id, customer_id, movie_id, showtime_id, amount, payment_method, seats, transaction_time)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (transaction_id) DO NOTHING
""", (
raw_data['transaction_id'],
customer['customer_id'],
movie['movie_id'],
showtime['showtime_id'],
payment['amount'],
payment['payment_method'],
raw_data['seats'],
raw_data['timestamp']
))
# Mark as processed
cursor.execute("""
UPDATE bronze_transactions SET processed = TRUE WHERE id = %s
""", (record_id,))
conn.commit()
cursor.close()
conn.close()
print(f"✅ Processed {len(records)} records from Bronze to Silver")
def refresh_gold_layer():
"""Refresh materialized view in Gold layer"""
pg_hook = PostgresHook(postgres_conn_id='postgres_cinema_dw')
conn = pg_hook.get_conn()
cursor = conn.cursor()
cursor.execute("REFRESH MATERIALIZED VIEW gold_cinema_analytics")
conn.commit()
cursor.close()
conn.close()
print("✅ Refreshed Gold layer materialized view")
# Define DAG
with DAG(
'bronze_to_silver_and_gold_elt',
default_args=default_args,
description='ELT pipeline: Bronze → Silver → Gold',
schedule_interval=timedelta(minutes=5),
catchup=False,
) as dag:
task_bronze_to_silver = PythonOperator(
task_id='bronze_to_silver',
python_callable=bronze_to_silver,
)
task_refresh_gold = PythonOperator(
task_id='refresh_gold_layer',
python_callable=refresh_gold_layer,
)
task_bronze_to_silver >> task_refresh_gold
```
### 5. Streamlit Dashboard
Real-time analytics visualization:
```python
# dashboard/app.py
import streamlit as st
import psycopg2
import pandas as pd
import plotly.express as px
import os
from time import sleep
st.set_page_config(page_title="CinéWorld Executive Dashboard", layout="wide")
@st.cache_resource
def get_db_connection():
return psycopg2.connect(
host=os.getenv('POSTGRES_HOST', 'localhost'),
port=os.getenv('POSTGRES_PORT', '5432'),
database=os.getenv('POSTGRES_DB', 'cinema_dw'),
user=os.getenv('POSTGRES_USER', 'postgres'),
password=os.getenv('POSTGRES_PASSWORD', 'postgres')
)
def fetch_gold_data():
"""Fetch aggregated analytics from Gold layer"""
conn = get_db_connection()
query = """
SELECT
cinema_location,
genre,
payment_method,
total_revenue,
ticket_count,
avg_ticket_price
FROM gold_cinema_analytics
"""
df = pd.read_sql(query, conn)
conn.close()
return df
# Dashboard header
st.title("🎬 CinéWorld Real-Time Analytics Dashboard")
st.markdown("**Live metrics from Gold layer** | Updates every 30 seconds")
# Auto-refresh
if 'refresh_counter' not in st.session_state:
st.session_state.refresh_counter = 0
placeholder = st.empty()
while True:
with placeholder.container():
# Fetch fresh data
df = fetch_gold_data()
# Metrics row
col1, col2, col3, col4 = st.columns(4)
with col1:
st.metric("Total Revenue", f"${df['total_revenue'].sum():,.2f}")
with col2:
st.metric("Total Tickets Sold", f"{df['ticket_count'].sum():,.0f}")
with col3:
GitHubで見る