Capstone: End-to-End Pipeline

Ten posts ago you couldn't draw the data platform. Now you're going to run one. This capstone wires every row of the platform diagram from post 1 into a single working pipeline: a stream of e-commerce order events flows through ingest → validate → transform → store → serve, and every stage runs for real on this lab machine — the event generator, the validator with its dead-letter queue, a Spark job with a genuine shuffle, a Parquet file with inspectable internals, and two serving systems answering the two questions they're built for. By the end you'll have watched the same numbers come out of three different systems and agree with each other. That agreement is the whole game.

Lab honesty, up front. Everything shown as a run below actually ran: the generator (5,000 events), the validator (113 rejects), the Spark RDD transform in local mode (real reduceByKey shuffle), the Parquet write/read with footer inspection, the Redis serving demo (Jedis), and the ClickHouse analytics (JDBC) — all on Temurin JDK 17, Spark 3.5.3, Redis 8.10.2, ClickHouse 26.10.1.1704. Three things did not run and are labeled [NOT RUN] where they appear: Spark's SQL/DataFrame engine (this sandbox's Netty transport fails on it — verified, repeatedly), Structured Streaming execution (same root cause), and Kafka as the pipe (not installed here; the generator stands in for it). The orchestration and monitoring sections are architectural guidance, labeled as such. Nothing below is presented as a run unless it says so.

The pipeline, all five rows

Here's the whole capstone in one diagram. Read it left to right — data only flows one way, and every arrow is a handoff you could observe:

ShopStream: one pipeline, five stages INGEST EventGen 5,000 events ran VALIDATE rules + dead-letter ran TRANSFORM Spark RDD reduceByKey ran STORE Parquet file footer inspected ran SERVE Redis + ClickHouse ran events.jsonl → valid.jsonl + deadletter.jsonl → revenue_csv → revenue_by_country.parquet → Redis hash + ClickHouse table Kafka would sit between ingest and validate in production [NOT RUN here — the generator stands in]

The scenario is ShopStream, a small e-commerce shop: every order emits an event {order_id, user_id, country, amount, ts}. The pipeline answers two questions at the end: "what's the revenue for US right now?" (Redis — a point lookup) and "how did revenue break down by country?" (ClickHouse — an analytical scan). Post 1's taxonomy, running.

Stage 1 — Ingest: generate the event stream

In production this stage is Kafka (the durable log from the Data & Messaging track). Here, a small Java generator plays the role of the event source: 5,000 order events with a fixed random seed, so every run produces the same stream. About 2% are deliberately invalid — negative amounts, unknown countries, empty user IDs — because a pipeline that never sees bad data is a pipeline that hasn't been tested. (In the commands below, $CP is the classpath: Jackson from Spark's jars for JSON parsing; the serving stages add Jedis and the ClickHouse JDBC driver.)

$ java -cp "$CP:." EventGen
GENERATED events=5000 invalid_injected=113 -> /tmp/capstone-lab/events.jsonl

$ head -1 events.jsonl
{"order_id":"ord-100000","user_id":"user-371","country":"IN","amount":343.2,"ts":1760000000}

Deterministic input is a testing decision, not laziness: when the input is fixed, every downstream number becomes a regression test. If the validator ever reports anything other than 113 rejects, something changed and you want to know.

Stage 2 — Validate: reject loudly, never silently

Validation is the cheapest correctness you can buy. Three rules — amount must be positive, country must be in the allowlist, user ID must be non-empty — and every reject lands in a dead-letter file with its reason attached. The dead-letter queue is the difference between "we lost 2% of events somewhere" and "we can show you exactly which 113 and why":

$ java -cp "$CP:." Validate
VALIDATED ok=4887 rejected=113 {amount_not_positive=41, empty_user_id=30, unknown_country=42}
  valid -> /tmp/capstone-lab/valid.jsonl
  dead-letter -> /tmp/capstone-lab/deadletter.jsonl

Note the reject counts add up: 41 + 30 + 42 = 113, and 4887 + 113 = 5000. Reconciliation is a habit, not a feature — at every stage boundary in this pipeline, the numbers in must equal the numbers out plus the numbers quarantined. The moment they don't, you stop and look.

Stage 3 — Transform: Spark does the shuffle

Now the compute stage: revenue and order count per country, from the 4,887 validated events. This runs on Spark's RDD API in local mode — and here's the pleasant surprise from the lab work: the real reduceByKey shuffle works fine in this sandbox. (The Netty transport failure here is specific to the SQL engine's execution path — the RDD shuffle path doesn't hit it.) The DataFrame equivalent of this job is [NOT RUN]: it can't execute here, so this post shows the RDD version that did.

$ spark-shell --master "local[2]" < transform.scala
  GB   revenue=   245206.71 orders=978
  DE   revenue=   190812.71 orders=762
  IN   revenue=   239076.62 orders=946
  BR   revenue=   127947.85 orders=508
  US   revenue=   428955.15 orders=1693

The script parses each JSON line with Jackson, maps to (country, (amount, 1)), and reduceByKey sums revenue and counts per country across a real shuffle. A second pass of the same script saved the five aggregate rows as CSV for the next stage:

$ spark-shell --master "local[2]" < transform_save.scala
SPARK_SAVED 5 rows

Sanity check before moving on: 978 + 762 + 946 + 508 + 1693 = 4887 orders. Every validated event is accounted for; nothing fell out of the shuffle.

Stage 4 — Store: the Parquet file, inspected

The transform's five aggregate rows become a Parquet file — written and read with pure-Java parquet-mr, no Spark engine involved. Small files like this don't show off Parquet's compression, but the footer inspection shows the columnar machinery is real: one row group, per-column encodings, exact row counts in metadata:

$ java -cp "$CP:." StoreParquet
PARQUET_WROTE rows=5 -> /tmp/capstone-lab/revenue_by_country.parquet
PARQUET_READ rows=5 revenue_sum=1231999.04 orders_sum=4887
PARQUET_FOOTER row_groups=1 rows_total=5
  column=[country] encodings=[BIT_PACKED, PLAIN] codec=UNCOMPRESSED
  column=[revenue] encodings=[BIT_PACKED, PLAIN] codec=UNCOMPRESSED
  column=[orders] encodings=[BIT_PACKED, PLAIN] codec=UNCOMPRESSED
PARQUET_SIZE_BYTES=755

Cross-check: revenue_sum 1,231,999.04 and 4,887 orders — identical to the Spark output. The Parquet file is the pipeline's durable record: if Redis and ClickHouse both vanished tonight, this file (or the lake it stands in for) rebuilds them.

Stage 5 — Serve: two systems, two questions

This is where post 1's taxonomy pays off. The same pipeline outputs serve two different access patterns, and forcing one system to do both jobs is the expensive mistake this whole track warns against.

Redis: "what's the revenue for US right now?"

Per-country revenue goes into a Redis hash; top users by order count go into a sorted set. Then the point lookups a live service would make — one hash field, one top-N:

$ java -cp "$CP:." ServeRedis
REDIS_PING -> PONG
REDIS_WROTE hash capstone:revenue:country fields=5
REDIS_WROTE zset capstone:leaderboard:orders_by_user members=400
REDIS_LOOKUP HGET capstone:revenue:country US -> 428955.15
REDIS_LOOKUP HGETALL -> {BR=127947.85, DE=190812.71, GB=245206.71, IN=239076.62, US=428955.15}
REDIS_LOOKUP TOP-3 users by orders -> [user-316, user-119, user-352] scores=[25.0, 25.0, 24.0]
REDIS_TTL capstone:pipeline:last_run -> 3600s

US revenue: 428,955.15 — the same number Spark computed. A product page or dashboard widget reads one hash field directly; it never touches the analytical store.

ClickHouse: "how did revenue break down by country?"

The 4,887 validated events land in a MergeTree table ordered by (country, ts), and the analytical query scans the amount column — not whole rows:

$ java -cp "$CP:." ServeClickHouse
CH_VERSION -> 26.10.1.1704
CH_INSERTED rows=4887 into capstone_orders
CH_ANALYTICS revenue by country:
  BR orders=508 revenue=127947.85
  DE orders=762 revenue=190812.71
  GB orders=978 revenue=245206.71
  IN orders=946 revenue=239076.62
  US orders=1693 revenue=428955.15

Three systems, one truth: Spark, Redis, and ClickHouse all report US = 428,955.15 across 1,693 orders. Principle: a pipeline earns trust one cross-check at a time. No single system's output is "correct" because it ran — it's correct because an independent system computed the same answer from the same input.

What production adds: orchestration and observability

The lab above runs the stages by hand, in order. Production doesn't — it schedules, retries, and watches. Two additions turn this pipeline from a demo into a system; both are architectural guidance here, not lab runs.

Orchestration: the DAG

A scheduler like Airflow or Dagster owns three things the shell script doesn't: ordering with dependencies (transform never runs before validation succeeds), retries with backoff (a transient ClickHouse timeout retries the load step, not the whole pipeline), and backfills (re-run last Tuesday's partition when a bug is found). The DAG for this pipeline is five tasks mirroring the five stages, each one idempotent — re-running the Spark transform for the same input partition must produce the same five rows, never double-count. Idempotency is what makes retries safe, and safe retries are what make a pipeline operable. [Architectural guidance — no scheduler was installed in this lab.]

Observability: the three signals, applied to the pipeline

The Production Java track taught logs, metrics, and traces for services; pipelines need the same three, aimed at data:

  • Logs (what happened to one run): every stage logs structured JSON with a shared run_id — events in, rejects with reasons, rows written. When the 4 AM alert fires, you grep one ID and see the whole story.
  • Metrics (what's happening across runs): events_in, events_rejected (by reason), revenue_sum as a gauge, stage durations. Alert on changes, not absolutes: a dead-letter spike from 2% to 12% means the source changed its schema; revenue drifting 30% from yesterday's same hour means something subtler broke.
  • Traces (where time went in one run): less critical for batch than for services, but stage-level timing (the transform stage dominates at roughly half a minute; ingest and serving loads take seconds) tells you where to optimize when the nightly window shrinks.

The single most valuable monitor for this pipeline is the reconciliation check from stage 2, automated: events in = events stored + events quarantined, every run, or page someone. [Architectural guidance — described, not run.]

The streaming version [NOT RUN]

Everything above is batch: bounded input, run to completion. The streaming version of ShopStream — Structured Streaming reading the event source continuously, watermarking late events, updating the Redis hash and ClickHouse table as micro-batches land — is the natural next step and the subject of post 6's concepts. It is [NOT RUN] here: Structured Streaming executes on the same SQL engine that can't run in this sandbox. On a normal machine the shape is: read stream → validate per micro-batch → foreachBatch writing to Redis and ClickHouse with idempotent sinks (the same idempotency that makes batch retries safe makes streaming retries safe). The exactly-once story from post 6 applies unchanged.

Cheat sheet

  • Five stages, one direction: ingest → validate → transform → store → serve. Data flows one way; every arrow is an observable handoff.
  • Deterministic input makes every downstream number a regression test: 5,000 in, 4,887 valid, 113 quarantined — with reasons.
  • Validate loudly: a dead-letter queue with reject reasons beats silent data loss every time.
  • Spark's RDD shuffle (reduceByKey) runs fine locally; the SQL/DataFrame engine is a separate beast with separate requirements.
  • Parquet is the durable record: the file's footer (row groups, encodings, row counts) is inspectable without any engine.
  • Redis answers "what about this one?" (point lookups); ClickHouse answers "how did it trend?" (columnar scans). One system doing both jobs badly is how you got here.
  • A pipeline earns trust one cross-check at a time. Spark, Redis, and ClickHouse agreeing on US = 428,955.15 is worth more than any single system's green checkmark.
  • Production adds orchestration (dependencies, retries, backfills — all resting on idempotency) and observability (logs per run, metrics across runs, reconciliation alerts).
  • Kafka would sit between ingest and validate in production; the generator stands in here. The exactly-once and watermarking story lives in post 6.

Field check before you close the track

(1) Re-run the whole pipeline and change the generator seed from 42 to 7 — which numbers change, which stay the same, and why? (2) Add a fourth validation rule (e.g. amount < 10,000 to catch fat-fingered entries) and confirm the dead-letter counts still reconcile. (3) The Redis hash and the ClickHouse table disagree after a re-run — walk through the three most likely causes in order. (4) Sketch the Airflow/Dagster DAG for these five stages: which task depends on which, what's idempotent, and what happens when the ClickHouse load fails twice?

Where to go from here

You started this track unable to draw the data platform. You can now do more than draw it: you can name the question each piece answers, choose storage by query pattern instead of hype, and — as of today — run an event through five stages and watch three systems agree on the answer. That's the job.

What this track didn't cover, deliberately: Flink as a primary engine (post 1 mapped it; a full Flink track would be its own series), table formats beyond the honest mention (Delta/Iceberg/Hudi deserve labs on a machine where the SQL engine runs), and the ML side of the platform (feature stores, training pipelines). The foundations you built here — partitions, shuffles, columnar storage, tunable consistency, serving trade-offs — transfer directly to all of them.

If you came from the Production Java track, the through-line is observability: the pipeline you just ran is another production system, and it deserves the same three signals. If you came from the Data & Messaging track, the through-line is the log: Kafka is the pipe this capstone's generator was standing in for. The platform diagram from post 1 is yours now — keep it somewhere visible.

Continue: Java Learning Roadmap 2026

Comments

Popular posts from this blog

JSP Servlet Interview Questions For Freshers Series 1

Java Banking Finance Services and Insurance (BFSI) domain interview questions

Java program to check even or odd number