A raw-to-gold lakehouse proof of concept on Cloudera
Two parallel pipelines, one for text and one for nested JSON, from generated events to partitioned Hive tables.
One event, raw to gold
Illustrative · synthetic dataA synthetic event, following the same steps the proof of concept runs.
Synthetic JSON event · E-0007
What changed here? · Generate
Before
Nothing yet. The record is created here.
Action
A generator writes events every minute, as nested JSON and as pipe-delimited text.
After
- {"event_id":"E-0007","type":"checkin","ts":"2026-03-18T09:15:02Z","site":{"region":"North","unit":4},"count":3}
What changed here? · Land
Before
- Local file
Action
Ingest adds a metadata header, lands the file in HDFS raw, archives the local copy and logs it as ingested.
After
- raw/json/2026-03-18/batch_0915.json
- ingested.log +1, so a rerun skips it
What changed here? · Transform
Before
- Nested JSON
Action
PySpark flattens nested fields and validates types.
After
- event_id=E-0007 · type=checkin · site_region=North · site_unit=4 · count=3
What changed here? · Gold table
Before
- Validated rows
Action
Rows are written as partitioned Parquet in Hive; reprocessing a day replaces only that partition.
After
- gold.events / event_date=2026-03-18
- SELECT count(*) FROM gold.events WHERE event_date='2026-03-18'
All steps as text
- Generate: A generator writes events every minute, as nested JSON and as pipe-delimited text. Result: {"event_id":"E-0007","type":"checkin","ts":"2026-03-18T09:15:02Z","site":{"region":"North","unit":4},"count":3}
- Land: Ingest adds a metadata header, lands the file in HDFS raw, archives the local copy and logs it as ingested. Result: raw/json/2026-03-18/batch_0915.json; ingested.log +1, so a rerun skips it
- Transform: PySpark flattens nested fields and validates types. Result: event_id=E-0007 · type=checkin · site_region=North · site_unit=4 · count=3
- Gold table: Rows are written as partitioned Parquet in Hive; reprocessing a day replaces only that partition. Result: gold.events / event_date=2026-03-18; SELECT count(*) FROM gold.events WHERE event_date='2026-03-18'
Major challenges
- 2 formats
Text and nested JSON
Two parallel pipelines parse 14 text fields and flatten 20 JSON fields into one table design.
- Idempotent
Safe reruns
Processed and ingested logs mean a rerun never duplicates data.
- Per-day
Cheap reprocessing
Partitioned Parquet with dynamic overwrite replaces only one day's partition.
- 3 nodes
Proof-of-concept cluster
Runs end to end on a three-node CDP 7.1.9 cluster with a runbook for reviewers.
My contribution
Designed and built the whole proof of concept, from generators and ingestion to transforms, table definitions, orchestration and runbook.
Tools used
- Cloudera CDP 7.1.9
- PySpark
- HDFS
- Hive
- Parquet
- Python
- Bash
- cron
Outcomes
- Two end-to-end pipelines, for text and for nested JSON, running on a three-node CDP cluster.
- Config-driven with no hard-coded paths, plus a runbook for operators.
- Synthetic data throughout.
ArchitectureAnonymised component diagram
Generate
- Text events (pipe-delimited)
- Nested JSON events
Land
- HDFS raw zoneIngest metadata and idempotency log
Transform
- PySpark parse (14 fields)
- PySpark flatten (20 fields)
Serve
- Hive gold tablesPartitioned Parquet
Read the flows as text
- Text events (pipe-delimited) → HDFS raw zone
- Nested JSON events → HDFS raw zone
- HDFS raw zone → PySpark parse (14 fields)
- HDFS raw zone → PySpark flatten (20 fields)
- PySpark parse (14 fields) → Hive gold tables
- PySpark flatten (20 fields) → Hive gold tables
Full case studyProblem, decisions, implementation, rollout
Problem and constraints
A prospective government client wanted to see unstructured and semi-structured data land in a governed lakehouse, and become queryable, on Cloudera.
- Proof-of-concept scope. A three-node CDP 7.1.9 cluster with synthetic data.
- Repeatable. Reviewers had to be able to rerun it, and runs had to be idempotent.
My role and the team’s
I designed and built the proof of concept: generators, ingestion, transforms, table definitions, orchestration, runbook and documentation.
Key decisions and trade-offs
- Config-driven. Paths and settings live in configuration files, so the same code runs on another cluster with no edits.
- Idempotent ingestion. Processed and ingested logs mean rerunning a step doesn’t duplicate data.
- Partitioned Parquet with dynamic partition overwrite. Reprocessing one day replaces only that day’s partition.
- Bash orchestrators over a scheduler. A proof of concept didn’t justify standing up a workflow engine, so a simple orchestrator with step timing and abort-on-failure was enough.
Implementation
- A generator writes events every minute (cron), in pipe-delimited text and nested JSON.
- Ingest adds a metadata header, lands files in HDFS raw, archives the local copy and records what was processed.
- PySpark parses 14 text fields or flattens 20 JSON fields (including nested ones), validates them and writes partitioned Parquet into Hive gold tables, then repairs partitions.
Testing and rollout
Processed and ingested logs make reruns safe, and the master orchestrator aborts on the first failed step. A runbook documents setup, operation and troubleshooting.
Verified outcomes
Both pipelines ran end to end on the proof-of-concept cluster.
Domains Data platforms · Processing & workflows · Streaming & ingestion
Code The repository names the prospective client, so it isn't linked. A sanitised demo version is planned.