Proof of conceptProof of concept via BBI.ai · Government agency2026

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.

EventsRawSparkGold partitions

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}

All steps as text
  1. 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}
  2. 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
  3. Transform: PySpark flattens nested fields and validates types. Result: event_id=E-0007 · type=checkin · site_region=North · site_unit=4 · count=3
  4. 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

  1. 2 formats

    Text and nested JSON

    Two parallel pipelines parse 14 text fields and flatten 20 JSON fields into one table design.

  2. Idempotent

    Safe reruns

    Processed and ingested logs mean a rerun never duplicates data.

  3. Per-day

    Cheap reprocessing

    Partitioned Parquet with dynamic overwrite replaces only one day's partition.

  4. 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
Both pipelines share the same three stages. An orchestrator times each step and stops on the first failure.
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

  1. A generator writes events every minute (cron), in pipe-delimited text and nested JSON.
  2. Ingest adds a metadata header, lands files in HDFS raw, archives the local copy and records what was processed.
  3. 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

Trace it in the constellation

Code The repository names the prospective client, so it isn't linked. A sanitised demo version is planned.

Search the portfolio

Try:

↑ ↓ to move · Enter to open · Esc to close