Data Engineering
Real-time Insurance Claims Streaming Pipeline
An end-to-end real-time streaming pipeline that ingests synthetic mobile money and insurance claim events (PaySim …
View Project
An end-to-end data pipeline that ingests Spotify streaming data via the Spotify API, stages it in Amazon S3, transforms it with AWS Glue (PySpark), loads it into Snowflake, and surfaces insights through a Power BI dashboard — covering top tracks, artist reach, and listening behaviour by time of day and geography.
import json
import boto3
import requests
from datetime import datetime, timezone
from base64 import b64encode
def get_spotify_token() -> str:
credentials = b64encode(
f"{SPOTIFY_CLIENT_ID}:{SPOTIFY_CLIENT_SECRET}".encode()
).decode()
response = requests.post(
"https://accounts.spotify.com/api/token",
headers={"Authorization": f"Basic {credentials}"},
data={"grant_type": "client_credentials"},
)
response.raise_for_status()
return response.json()["access_token"]
def fetch_recently_played(token: str, limit: int = 50) -> list[dict]:
response = requests.get(
"https://api.spotify.com/v1/me/player/recently-played",
headers={"Authorization": f"Bearer {token}"},
params={"limit": limit},
)
response.raise_for_status()
return response.json().get("items", [])
def upload_to_s3(data: list[dict], bucket: str, prefix: str) -> str:
s3 = boto3.client("s3")
now = datetime.now(timezone.utc)
key = (
f"{prefix}/year={now.year}/month={now.month:02d}/"
f"day={now.day:02d}/{now.strftime('%H%M%S')}.json"
)
s3.put_object(
Bucket=bucket,
Key=key,
Body=json.dumps(data, ensure_ascii=False),
ContentType="application/json",
)
return f"s3://{bucket}/{key}"
def run_ingestion():
token = get_spotify_token()
tracks = fetch_recently_played(token)
s3_path = upload_to_s3(tracks, S3_BUCKET, S3_PREFIX)
print(f"Ingested {len(tracks)} records → {s3_path}")
if __name__ == "__main__":
run_ingestion()
import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, BooleanType
args = getResolvedOptions(
sys.argv, ["JOB_NAME", "source_path", "target_path"]
)
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args["JOB_NAME"], args)
# --- Read raw JSON from Bronze layer ---
df_raw = (
spark.read
.option("multiline", True)
.json(args["source_path"])
)
# --- Flatten nested track / artist / context fields ---
df_flat = df_raw.select(
F.col("played_at").cast("timestamp").alias("played_at"),
F.col("track.id").alias("track_id"),
F.col("track.name").alias("track_name"),
F.col("track.duration_ms").cast(LongType()).alias("duration_ms"),
F.col("track.popularity").alias("popularity"),
F.col("track.explicit").cast(BooleanType()).alias("is_explicit"),
F.col("track.artists")[0]["id"].alias("primary_artist_id"),
F.col("track.artists")[0]["name"].alias("primary_artist_name"),
F.col("track.album.id").alias("album_id"),
F.col("track.album.name").alias("album_name"),
F.col("track.album.release_date").alias("release_date"),
F.col("context.type").alias("context_type"),
F.col("context.uri").alias("context_uri"),
)
# --- Deduplicate on natural key ---
df_deduped = df_flat.dropDuplicates(["played_at", "track_id"])
# --- Add date partition columns ---
df_final = (
df_deduped
.withColumn("year", F.year("played_at"))
.withColumn("month", F.month("played_at"))
.withColumn("day", F.dayofmonth("played_at"))
.withColumn("hour", F.hour("played_at"))
)
# --- Data quality: drop rows missing critical keys ---
df_final = df_final.filter(
F.col("track_id").isNotNull() & F.col("played_at").isNotNull()
)
print(f"Writing {df_final.count()} clean records to Silver layer...")
# --- Write to Silver layer as partitioned Parquet ---
(
df_final.write
.mode("overwrite")
.partitionBy("year", "month", "day")
.parquet(args["target_path"])
)
job.commit()
print("Glue job complete.")
-- ============================================================
-- 1. SNOWFLAKE SETUP
-- ============================================================
CREATE WAREHOUSE IF NOT EXISTS spotify_wh
WAREHOUSE_SIZE = 'X-SMALL'
AUTO_SUSPEND = 60
AUTO_RESUME = TRUE;
CREATE DATABASE IF NOT EXISTS spotify_db;
CREATE SCHEMA IF NOT EXISTS spotify_db.gold;
USE SCHEMA spotify_db.gold;
-- ============================================================
-- 2. DIMENSION: dim_tracks
-- ============================================================
CREATE OR REPLACE TABLE dim_tracks (
track_id VARCHAR(50) NOT NULL PRIMARY KEY,
track_name VARCHAR(500),
album_id VARCHAR(50),
album_name VARCHAR(500),
release_date DATE,
duration_ms BIGINT,
is_explicit BOOLEAN,
popularity INTEGER,
loaded_at TIMESTAMP_NTZ DEFAULT CURRENT_TIMESTAMP()
);
-- ============================================================
-- 3. DIMENSION: dim_artists
-- ============================================================
CREATE OR REPLACE TABLE dim_artists (
artist_id VARCHAR(50) NOT NULL PRIMARY KEY,
artist_name VARCHAR(300),
loaded_at TIMESTAMP_NTZ DEFAULT CURRENT_TIMESTAMP()
);
-- ============================================================
-- 4. FACT: fact_streams
-- ============================================================
CREATE OR REPLACE TABLE fact_streams (
stream_id VARCHAR(100) NOT NULL PRIMARY KEY,
played_at TIMESTAMP_NTZ NOT NULL,
track_id VARCHAR(50) REFERENCES dim_tracks(track_id),
primary_artist_id VARCHAR(50) REFERENCES dim_artists(artist_id),
context_type VARCHAR(50), -- e.g. playlist, album, artist
context_uri VARCHAR(500),
year INTEGER,
month INTEGER,
day INTEGER,
hour INTEGER,
loaded_at TIMESTAMP_NTZ DEFAULT CURRENT_TIMESTAMP()
)
CLUSTER BY (year, month, day); -- Partition pruning for date-range queries
-- ============================================================
-- 5. DBT GOLD MODEL: daily_stream_summary.sql
-- Materialised as a table, refreshed daily
-- ============================================================
-- models/gold/daily_stream_summary.sql
-- {{ config(materialized='table') }}
SELECT
DATE_TRUNC('day', s.played_at) AS stream_date,
t.track_name,
a.artist_name,
s.context_type,
COUNT(*) AS total_streams,
SUM(t.duration_ms) / 1000 / 60.0 AS total_minutes_played,
AVG(t.popularity) AS avg_popularity_score,
COUNT(DISTINCT s.context_uri) AS unique_contexts
FROM fact_streams s
JOIN dim_tracks t ON s.track_id = t.track_id
JOIN dim_artists a ON s.primary_artist_id = a.artist_id
GROUP BY 1, 2, 3, 4
ORDER BY stream_date DESC, total_streams DESC;
-- ============================================================
-- 6. POWER BI CONNECTOR QUERY
-- Power BI > Get Data > Snowflake
-- ============================================================
SELECT
stream_date,
track_name,
artist_name,
context_type,
total_streams,
ROUND(total_minutes_played, 2) AS total_minutes_played,
ROUND(avg_popularity_score, 1) AS avg_popularity_score,
unique_contexts
FROM spotify_db.gold.daily_stream_summary
WHERE stream_date >= DATEADD('day', -90, CURRENT_DATE())
ORDER BY stream_date DESC;
An end-to-end real-time streaming pipeline that ingests synthetic mobile money and insurance claim events (PaySim …
View Project