This is a local offline advertising analytics warehouse built with Hive Metastore + Spark SQL + PySpark + Parquet.
This project models a realistic ad-data pipeline from raw event ingestion to reporting outputs through a classic ODS -> DWD -> DWS -> ADS architecture. It includes fact and dimension ingestion, SQL-first warehouse jobs, a user-tagging module, and a reproducible benchmark for data skew and salting-based optimization.
- Four-layer warehouse architecture:
ODS -> DWD -> DWS -> ADS - Hive Metastore used as the metadata/catalog layer, with Parquet as the physical storage layer
- SQL-first DWD / DWS / ADS jobs, with Python kept as a lightweight orchestrator
- Explicit fact and dimension modeling for event, conversion, cost, user, ad, and ad-slot data
- User tag snapshot pipeline based on daily and rolling 7-day features
- ADS outputs for KPI overview, campaign ranking, advertiser dashboard, and tag effectiveness
- Reproducible 20M-row benchmark for campaign-level data skew
- Salting + two-stage aggregation experiment to mitigate hot-key skew on
campaign_id
- Python
- PySpark
- Spark SQL
- Hive Metastore (embedded Derby for local development)
- Parquet
- Local Spark (
local[*]) - Benchmark / skew experiment tooling under
benchmark/
Raw CSV generator
↓
ODS
↓
DWD
↓
DWS
↓
ADS
- Ingests raw fact CSVs into partitioned Parquet under
warehouse/ods/... - Registers ODS facts and dimensions into the
odsHive database - Stores:
- event log
- conversion log
- campaign cost
- user profile
- ad metadata
- ad slot metadata
- Cleans and standardizes raw ODS inputs
- Deduplicates facts, filters invalid traffic, and normalizes event / conversion timestamps
- Enriches facts with business dimensions such as advertiser, ad type, landing type, product, and slot metadata
- Produces dimension snapshots and detail facts:
- user dimension
- ad dimension
- ad-slot dimension
- impression detail
- click detail
- conversion detail
- campaign-day cost
- Builds daily subject-area aggregates on top of DWD
- Standardizes campaign, advertiser, and user metrics such as
impressions,clicks,conversions,gmv,cost,ctr,cvr,rpm, androi - Adds a user-tagging module:
- static tag dictionary
- daily user tag snapshot
- daily tag quality table
- Produces reporting-friendly outputs from DWS
- Includes:
- KPI overview
- campaign ranking
- advertiser dashboard
- tag effectiveness report
| Table | Grain | Purpose |
|---|---|---|
ods.ad_event_log |
dt + event_id |
Raw impression / click event stream landed to Parquet and Hive |
ods.conversion_log |
dt + conv_id |
Conversion facts with order and GMV information |
ods.ad_cost |
dt + campaign_id |
Daily campaign cost input for downstream ROI/cost analysis |
ods.user_profile |
user_id |
User dimension source with demographics and registration date |
ods.ad_meta |
ad_id + campaign_id |
Ad metadata source with advertiser, ad type, landing type, and product |
ods.ad_slot |
ad_slot_id |
Ad slot / placement metadata |
| Table | Grain | Purpose |
|---|---|---|
dwd.dim_user |
dt + user_id |
Daily user dimension snapshot |
dwd.dim_ad |
dt + ad_id + campaign_id |
Daily ad dimension snapshot aligned to fact enrichment grain |
dwd.dim_ad_slot |
dt + ad_slot_id |
Daily ad-slot dimension snapshot |
dwd.ad_impression_detail |
dt + event_id |
Cleaned and enriched impression facts |
dwd.ad_click_detail |
dt + event_id |
Cleaned and enriched click facts |
dwd.ad_conversion_detail |
dt + conv_id |
Cleaned and enriched conversion facts |
dwd.campaign_day_cost |
dt + campaign_id |
Standardized campaign-day cost table for DWS cost aggregation |
| Table | Grain | Purpose |
|---|---|---|
dws.campaign_daily |
dt + campaign_id |
Daily campaign metrics for performance analysis |
dws.advertiser_daily |
dt + advertiser_id |
Daily advertiser aggregates for dashboarding |
dws.user_daily |
dt + user_id |
Daily user behavior aggregate |
dws.user_tag_snapshot |
dt + user_id |
User tag snapshot built from same-day and 7-day rolling features |
dws.tag_quality_daily |
dt + tag_name |
Daily tag coverage and effectiveness metrics |
dws.tag_dict |
tag_id |
Built-in tag dictionary and rule definitions |
| Table | Grain | Purpose |
|---|---|---|
ads.kpi_overview_daily |
dt |
Executive KPI overview of traffic, conversion, GMV, cost, and ROI |
ads.campaign_ranking_daily |
dt + campaign_id |
Campaign leaderboard by GMV and ROI |
ads.advertiser_dashboard_daily |
dt + advertiser_id |
Advertiser dashboard with growth and retention indicators |
ads.tag_effectiveness_daily |
dt + tag_name |
Tag-level effectiveness summary |
jobs/ingest_ods.pylands ODS fact tables to partitioned Parquet and registers them in Hive.jobs/ingest_ods_dims.pyingests user / ad / ad-slot dimensions into the ODS layer.jobs/build_dwd.pyis SQL-dominant and converts ODS data into cleaned, enriched warehouse facts and dimension snapshots.
- Implemented inside
jobs/build_dws.py - Uses both daily features and rolling 7-day features
- Produces:
dws.tag_dictdws.user_tag_snapshotdws.tag_quality_daily
jobs/build_ads.pygenerates dashboard-ready outputs from DWS- Covers overview KPIs, campaign ranking, advertiser reporting, and tag effectiveness analysis
The project includes a standalone benchmark under benchmark/ to reproduce and explain data skew on campaign_id aggregation.
Validate how a hot campaign_id affects Spark aggregation runtime, partition balance, and resource pressure under a campaign-level aggregation workload similar to DWS.
- Dataset size:
20,000,000rows per distribution - Hot-key setup: one
campaign_idaccounts for roughly80%of rows in the skewed dataset - Stress benchmark mode:
- fixed shuffle partitions
- repartition by
campaign_id - sort within partitions
- campaign-level aggregation
Observed locally with 20M rows and shuffle_partitions=4:
| Scenario | Elapsed Time | Partition Skew Ratio |
|---|---|---|
uniform |
36.31s |
1.14 |
skewed |
62.67s |
3.43 |
slowdown_ratio = 1.73- The skewed run also triggered memory-pressure warnings during execution
The benchmark also includes a salted version of the skewed aggregation:
- Input is still the same
skeweddataset - The hot
campaign_idis split into multiple salted sub-keys - Aggregation is executed in two stages:
dt + campaign_id + saltdt + campaign_id
Observed locally with 20M rows, shuffle_partitions=4, and salt_buckets=8:
| Scenario | Elapsed Time | Partition Skew Ratio |
|---|---|---|
skewed |
32.03s |
3.43 |
skewed_salted |
23.95s |
1.80 |
improvement_ratio = 1.34- The salted run materially reduced partition imbalance and improved runtime
- A hot
campaign_idcan create strong partition imbalance, long-tail tasks, and higher resource pressure during aggregation. - Salting + two-stage aggregation is an effective mitigation strategy for hot-key skew, even in a local Spark environment.
python -m venv .venv
source .venv/bin/activate
pip install -r requirements.txtFor a realistic DWS tagging run, generate at least 7 days:
python data_generator/generate_ods.py --start_dt 2026-03-01 --days 7PYTHONPATH=. python jobs/init_hive.pyPYTHONPATH=. python jobs/ingest_ods.py
PYTHONPATH=. python jobs/ingest_ods_dims.pyIf you want correct 7-day user-tag features, run DWD and DWS sequentially by day:
for dt in 2026-03-01 2026-03-02 2026-03-03 2026-03-04 2026-03-05 2026-03-06 2026-03-07; do
PYTHONPATH=. python -m jobs.build_dwd --dt "$dt" --no_dq
PYTHONPATH=. python -m jobs.build_dws --dt "$dt" --no_dq
donePYTHONPATH=. python -m jobs.build_ads --dt 2026-03-07 --no_dqFor a quick end-to-end smoke run:
bash run_all.shNote: run_all.sh is useful for a single-date smoke test, while the 7-day sequential DWD/DWS run above is recommended when demonstrating rolling user-tag features.
python benchmark/generate_skew_data.py --rows 20000000PYTHONPATH=. python benchmark/run_campaign_skew_benchmark.py --modes all --benchmark_mode stress --shuffle_partitions 4 --salt_buckets 8PYTHONPATH=. python benchmark/run_campaign_skew_benchmark.py --modes skewed skewed_salted --benchmark_mode stress --shuffle_partitions 4 --salt_buckets 8- Data:
benchmark/data/... - Aggregation outputs:
benchmark/results/<mode>/campaign_daily_stress - Result files:
benchmark/results/benchmark_results.csvbenchmark/results/benchmark_results.jsonl
.
├── benchmark/
│ ├── generate_skew_data.py
│ ├── run_campaign_skew_benchmark.py
│ ├── data/
│ ├── results/
│ └── tmp/
├── common/
│ └── spark_session.py
├── data/
│ └── ods/
├── data_generator/
│ └── generate_ods.py
├── docs/
├── jobs/
│ ├── init_hive.py
│ ├── ingest_ods.py
│ ├── ingest_ods_dims.py
│ ├── build_dwd.py
│ ├── build_dws.py
│ ├── build_ads.py
│ └── dq_check.py
├── scripts/
├── warehouse/
│ ├── ods/
│ ├── dwd/
│ ├── dws/
│ ├── ads/
│ ├── hive_warehouse/
│ └── metastore_db/
├── requirements.txt
├── run_all.sh
└── README.md