| 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 — 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
git clone https://github.com/BaidaneAyoub/realtime-cinema-data-engineering.git
cd realtime-cinema-data-engineering
python -m venv myenv
source myenv/bin/activate
pip install -r requirements.txt
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
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)
CREATE TABLE silver_customers (...);
CREATE TABLE silver_movies (...);
CREATE TABLE silver_showtimes (...);
CREATE TABLE silver_transactions (...);
Gold Layer: Materialized views for analytics
CREATE MATERIALIZED VIEW gold_cinema_analytics AS
SELECT ...
2. Kafka Event Producer
Generate and stream synthetic cinema transaction events:
from kafka import KafkaProducer
from faker import Faker
import json
import time
import os
fake = Faker()
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": {
: (fake.random.uniform(, ), ),
: fake.random_element([, , ]),
:
},
: [],
: fake.date_time_this_month().isoformat()
}
():
()
:
:
event = generate_ticket_sale()
producer.send(topic, value=event)
()
time.sleep(interval)
KeyboardInterrupt:
()
:
producer.flush()
producer.close()
__name__ == :
start_streaming()
3. Kafka Consumer (Bronze Layer Ingestion)
Consume events from Kafka and insert into PostgreSQL Bronze layer:
from kafka import KafkaConsumer
import psycopg2
from psycopg2.extras import Json
import json
import os
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'))
)
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()
def start_consuming():
()
conn = get_db_connection()
:
message consumer:
event = message.value
insert_bronze(conn, event)
()
KeyboardInterrupt:
()
:
conn.close()
consumer.close()
__name__ == :
start_consuming()
4. Airflow ELT DAG (Bronze → Silver → Gold)
Orchestrate data transformation pipeline:
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()
cursor.execute("""
SELECT id, raw_data
FROM bronze_transactions
WHERE processed = FALSE
LIMIT 50000
""")
records = cursor.fetchall()
for record_id, raw_data in records:
customer = raw_data['customer']
movie = raw_data['movie']
showtime = raw_data['showtime']
payment = raw_data['payment']
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)
cursor.execute(, movie)
cursor.execute(, showtime)
cursor.execute(, (
raw_data[],
customer[],
movie[],
showtime[],
payment[],
payment[],
raw_data[],
raw_data[]
))
cursor.execute(, (record_id,))
conn.commit()
cursor.close()
conn.close()
()
():
pg_hook = PostgresHook(postgres_conn_id=)
conn = pg_hook.get_conn()
cursor = conn.cursor()
cursor.execute()
conn.commit()
cursor.close()
conn.close()
()
DAG(
,
default_args=default_args,
description=,
schedule_interval=timedelta(minutes=),
catchup=,
) dag:
task_bronze_to_silver = PythonOperator(
task_id=,
python_callable=bronze_to_silver,
)
task_refresh_gold = PythonOperator(
task_id=,
python_callable=refresh_gold_layer,
)
task_bronze_to_silver >> task_refresh_gold
5. Streamlit Dashboard
Real-time analytics visualization:
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
st.title("🎬 CinéWorld Real-Time Analytics Dashboard")
st.markdown("**Live metrics from Gold layer** | Updates every 30 seconds")
if 'refresh_counter' not in st.session_state:
st.session_state.refresh_counter = 0
placeholder = st.empty()
:
placeholder.container():
df = fetch_gold_data()
col1, col2, col3, col4 = st.columns()
col1:
st.metric(, )
col2:
st.metric(, )
col3:
st.metric(, )
col4:
st.metric(, df[].nunique())
col_left, col_right = st.columns()
col_left:
location_revenue = df.groupby()[].().reset_index()
fig1 = px.bar(
location_revenue,
x=,
y=,
title=,
labels={: , : }
)
st.plotly_chart(fig1, use_container_width=)
payment_tickets = df.groupby()[].().reset_index()
fig3 = px.pie(
payment_tickets,
values=,
names=,
title=
)
st.plotly_chart(fig3, use_container_width=)
col_right:
genre_revenue = df.groupby()[].().reset_index()
fig2 = px.treemap(
genre_revenue,
path=[],
values=,
title=
)
st.plotly_chart(fig2, use_container_width=)
st.subheader()
top_locations = df.groupby().agg({
: ,
:
}).sort_values(, ascending=).head()
st.dataframe(top_locations, use_container_width=)
sleep()
st.session_state.refresh_counter +=
st.rerun()
Configuration
Environment Variables
Create a .env file for configuration:
KAFKA_BOOTSTRAP_SERVERS=localhost:9092
KAFKA_TOPIC=cinema_transactions
POSTGRES_HOST=localhost
POSTGRES_PORT=5432
POSTGRES_DB=cinema_dw
POSTGRES_USER=postgres
POSTGRES_PASSWORD=postgres
AIRFLOW__CORE__EXECUTOR=LocalExecutor
AIRFLOW__CORE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow
AIRFLOW_CONN_POSTGRES_CINEMA_DW=postgresql://postgres:postgres@postgres:5432/cinema_dw
Docker Compose Services
services:
kafka:
image: confluentinc/cp-kafka:latest
ports:
- "9092:9092"
postgres:
image: postgres:15
environment:
POSTGRES_USER: postgres
POSTGRES_PASSWORD: postgres
POSTGRES_DB: cinema_dw
ports:
- "5432:5432"
airflow-webserver:
image: apache/airflow:2.7.0
ports:
- "8080:8080"
Common Workflows
Starting the Complete Pipeline
docker-compose up -d
sleep 180
source myenv/bin/activate
python producer/main_producer.py
source myenv/bin/activate
python consumer/main_consumer.py
source myenv/bin/activate
streamlit run dashboard/app.py
Manual DAG Trigger
docker exec -it <airflow-container-id> airflow dags trigger bronze_to_silver_and_gold_elt
Query Gold Layer Directly
import psycopg2
import pandas as pd
conn = psycopg2.connect(
host="localhost",
database="cinema_dw",
user="postgres",
password="postgres"
)
query = """
SELECT genre, SUM(total_revenue) as revenue
FROM gold_cinema_analytics
GROUP BY genre
ORDER BY revenue DESC
LIMIT 5
"""
df = pd.read_sql(query, conn)
print(df)
conn.close()
Troubleshooting
Kafka Consumer Not Receiving Messages
docker exec -it <kafka-container> kafka-topics --list --bootstrap-server localhost:9092
docker exec -it <kafka-container> kafka-consumer-groups \
--bootstrap-server localhost:9092 \
--group cinema-consumer-group \
--describe
Airflow DAG Not Running
from airflow.models import DagBag
dagbag = DagBag()
dag = dagbag.get_dag('bronze_to_silver_and_gold_elt')
print(f"DAG errors: {dagbag.import_errors}")
PostgreSQL Connection Issues
docker exec -it <postgres-container> psql -U postgres -d cinema_dw -c "\dt"
docker exec -it <postgres-container> psql -U postgres -d cinema_dw -c "
SELECT schemaname, tablename
FROM pg_tables
WHERE schemaname = 'public';
"
Streamlit Dashboard Not Updating
conn = psycopg2.connect(...)
cursor = conn.cursor()
cursor.execute("SELECT COUNT(*) FROM gold_cinema_analytics")
count = cursor.fetchone()[0]
print(f"Gold layer records: {count}")
Performance Optimization
Batch Processing Size
Adjust batch size in Airflow DAG:
cursor.execute("""
SELECT id, raw_data
FROM bronze_transactions
WHERE processed = FALSE
LIMIT 100000 -- Increase for better throughput
""")
Kafka Consumer Parallelization
consumer = KafkaConsumer(
'cinema_transactions',
bootstrap_servers='localhost:9092',
group_id='cinema-consumer-group',
max_poll_records=500,
session_timeout_ms=30000
)
Materialized View Incremental Refresh
CREATE MATERIALIZED VIEW gold_cinema_analytics AS
SELECT ...
WITH DATA;
CREATE UNIQUE INDEX ON gold_cinema_analytics (cinema_location, genre, payment_method);
REFRESH MATERIALIZED VIEW CONCURRENTLY gold_cinema_analytics;
This skill enables AI agents to help developers build production-grade real-time streaming data pipelines with proper architecture patterns, orchestration, and visualization.