Data Engineering on the JVM: The Landscape

Your Spring Boot service from the Data & Messaging track is doing well. Too well: the Postgres holding its orders just crossed 2 TB, the nightly revenue report — one big SQL query — now takes four hours and times out half the time, and product just asked for "a dashboard that updates every minute." Someone in the meeting says "we need a data platform." Everyone nods. Nobody can draw what that means.

This post draws the map. Data engineering on the JVM is a crowded territory — Hadoop, Spark, Flink, Cassandra, HBase, Redis, ClickHouse, Kafka — and every tool's marketing claims it does everything. It doesn't. Each one was built for one shape of work, and choosing wrong is the most expensive mistake in this space: a batch engine forced into real-time, or a serving store forced into analytics, fails slowly and expensively. By the end of this post you'll know what each piece does, why the JVM ended up running nearly all of it, and which tool fits which problem. Then you'll prove your own lab works with a real Spark run — because every post in this track backs its claims with commands that actually ran.

Lab honesty, up front. This post is a map, not a lab: its job is orientation. The one thing that ran here is the Spark smoke test at the end (Spark 3.5.3, local mode, Temurin JDK 17) — real commands, real output. The architecture descriptions (HDFS, YARN, Flink, Cassandra internals and the rest) are standard, documented knowledge, not lab runs; each gets its full hands-on treatment in its own post later in the track. Nothing below is presented as a run unless it says so.

Notice the smell in that opening scene: one Postgres is being asked to do two fundamentally different jobs — record orders reliably (transactional work) and scan terabytes for a revenue report (analytical work). The industry names for those jobs are OLTP (online transaction processing: many small reads and writes, where consistency matters) and OLAP (online analytical processing: few huge scans, where throughput matters). The first architectural lesson of this track: when one database is asked to do both badly, the answer isn't a bigger database — it's a data platform.

The three shapes of data work

Every data tool answers one of three questions. Confusing them is the root of most bad architecture:

Three shapes of data work BATCH "what happened?" bounded input, whole dataset latency: minutes to hours optimizes: throughput nightly revenue report, retrain the model weekly STREAMING "what's happening?" unbounded input, never ends latency: seconds optimizes: freshness fraud alerts, live dashboard, per-minute aggregates SERVING "what about this one?" point lookups, tiny reads latency: milliseconds optimizes: tail latency user profile at login, product page in 50 ms One question per workload. The tool that answers one well usually answers the others badly.

Batch asks "what happened?" — the input is bounded (yesterday's logs, the whole orders table), and you trade latency for throughput: scanning terabytes is fine if it finishes by morning. Streaming asks "what's happening?" — the input never ends, and you trade completeness for freshness: an approximate fraud score in 5 seconds beats a perfect one tomorrow. Serving asks "what about this one?" — tiny reads, millions of them, where the 99th percentile latency is the product.

Hold those three apart and the tool zoo organizes itself: Spark and MapReduce primarily handle batch processing. Flink and Spark Structured Streaming handle streaming processing. Cassandra, HBase, Redis, and ClickHouse sit downstream as storage or serving systems — each optimized for a very different kind of access, which is exactly where people go wrong. Kafka, from the Data & Messaging track, is the pipe between them: the log that streaming engines read and batch engines catch up on.

The Hadoop ecosystem: what each piece actually does

Hadoop is the old growth forest everything else grew out of. You will probably never write a MapReduce job by hand in production — but you need to know what each piece does, because Spark, Flink, and the lakehouse are all reactions to Hadoop's answers:

  • HDFS (storage). One filesystem spread across machines. Files are chopped into blocks (typically 128 or 256 MB), each block replicated — usually 3 copies on different machines and racks. One NameNode holds the metadata (which blocks make up which file); DataNodes hold the blocks. It's write-once-read-many: great for ingesting huge immutable datasets, terrible for updating a single record. If you've ever wondered why data lakes store Parquet files instead of rows in a database — this is why.
  • YARN (resource management). The cluster's landlord. It doesn't know what a "query" is; it hands out containers — a slice of CPU and RAM on some machine — to whatever framework asks. MapReduce, Spark, and Flink all learned to run on YARN instead of each demanding their own cluster. (Today Kubernetes is eating this role, but the idea is identical: separate resource scheduling from the compute framework.)
  • MapReduce (the original engine). Google's 2004 idea, open-sourced as Hadoop: express the job as a map function (record → key/value pairs) and a reduce function (key + values → result), and the framework handles distribution, retries, and the brutal middle step — the shuffle, where all values for one key are routed to one machine. Its fatal flaw for modern work: every stage writes intermediate results back to disk. Iterative algorithms (machine learning, graph processing) pay that disk cost on every pass.

The honest summary: Hadoop solved reliable storage and scheduling for commodity hardware and its answers are still the foundation. What it got wrong was the execution engine — and that's the opening Spark walked through.

Why Spark won

Spark (Berkeley AMPLab, 2009) kept Hadoop's good ideas — distributed datasets, fault tolerance through lineage instead of replication — and fixed the engine: it avoided MapReduce's requirement to materialize every intermediate stage to distributed storage. Its DAG execution model pipelines transformations, reuses explicitly cached data across iterative computations, and spills to disk when memory runs out. That dramatically reduces I/O for many workloads — especially iterative ones like machine learning, which paid MapReduce's disk cost on every pass — and it made one engine credibly cover batch, SQL, streaming, and ML instead of four separate systems.

Two design decisions explain its staying power. First, the DataFrame API: instead of opaque JVM objects, you work with named, typed columns — which lets the Catalyst optimizer rewrite your query plan (predicate pushdown, column pruning) the way a database would. You write df.filter("status = 'shipped'") and Spark figures out how to skip the irrelevant data. Second, it meets Java developers where they are: the Java API is a first-class citizen (lambdas, Encoders), it runs on YARN or Kubernetes or your laptop, and the same code scales from a CSV on disk to a hundred-node cluster. That's why this track's labs run real Spark, and why post 3 goes deep on the Java API specifically.

Where Flink fits

If Spark is a batch engine that learned streaming, Flink is a streaming engine that learned batch. The distinction matters when latency and correctness get serious: Flink was built around event time (when the event happened, not when it arrived), watermarks (how long to wait for late data), and exactly-once state via distributed snapshots — as first principles, not retrofits. For "process every payment exactly once, in event-time order, with 100 ms latency," Flink is the purpose-built tool; Spark Structured Streaming (post 6) gets you most of the way there with a simpler model.

Rule of thumb, not a law: if seconds-to-minutes latency is acceptable and Spark is already your platform, Structured Streaming is often enough; when consistently low latency, sophisticated event-time processing, and deeply stateful streaming are central requirements, Flink deserves serious consideration. Don't introduce a second streaming platform just because one benchmark says it's faster — most companies never outgrow the first half of that sentence.

Why the JVM runs the data world

Step back and notice something odd: Hadoop, Spark, Flink, Kafka, Elasticsearch, Cassandra — nearly the entire data-infrastructure canon is written in Java or Scala. That's not an accident, and it's not fashion. The JVM's success here is a combination of several reinforcing advantages:

  1. A portable runtime with mature concurrency. Write once, run on any cluster OS, and the Java Memory Model gives library authors precise, portable concurrency semantics — the foundation every one of these systems' threading models is built on.
  2. The JIT. Hot loops get compiled to near-native code, so the "slow language" objection rarely survives contact with a warmed-up data engine. (Many engines also go off-heap or use native memory precisely to sidestep GC pressure — the runtime gives you both options.)
  3. Garbage collectors tuned for allocation-heavy workloads. A shuffle stage can allocate gigabytes per second; G1 — and now ZGC/Shenandoah with single-digit-millisecond pauses — are the product of twenty-five years of exactly this workload. GC is part of the story, not the whole story.
  4. The library network effect. Once Hadoop was Java, everything that wanted to read HDFS had a reason to be JVM-native; once Spark was Scala, its ecosystem compounded. Data infrastructure is glue-heavy work, and the glue is all on Maven Central.
  5. Operational maturity. Twenty years of jstack, heap dumps, JMX, and profilers mean a JVM data system is debuggable in production — which, as the Production Java track argued, is where systems actually live.

None of this means Python is absent — it's the language of notebooks, orchestration glue, and ML modeling. But the engines themselves, the things moving the bytes at 3 AM, are JVM programs. That's why this track exists on this site: you already know the runtime; now learn what it was built to do at scale.

Downstream storage and serving: one honest comparison

Batch and streaming compute; downstream systems have to store and serve the results — and "NoSQL" is four different answers to four different access patterns. Here's the whole landscape in one table; posts 7–9 give each its lab:

SystemData modelAnswers which question?Know this about it
CassandraWide-column (partition key + clustering columns)"Give me the last 50 events for user X" — high-write, key-scoped readsDynamo-style: masterless, LSM-tree storage, tunable consistency through read/write consistency levels — it leans toward availability under partitions, but the consistency behavior is yours to choose. Built for write throughput across data centers. CQL looks like SQL and isn't — no joins, model the query first.
HBaseWide-column on HDFS (BigTable clone)"Give me row Y as of Tuesday" — strong per-row consistency on huge tablesFavors strong row-level consistency: one region server owns each row range, coordinated infrastructure underneath, living on your HDFS. You'll often hear "Cassandra leans AP, HBase leans CP" — treat that as architectural shorthand, not the complete consistency model. The choice when you already run Hadoop and need random reads with real consistency.
RedisIn-memory key/value + data structures"What's in user X's session right now?" — sub-millisecond servingMost commonly the speed layer — cache, session store, rate limiter, leaderboard — where memory-speed access matters. RDB snapshots and AOF persistence provide durability options, but using Redis as the primary system of record requires deliberate durability, replication, and recovery design.
ClickHouseColumnar OLAP (MergeTree engine)"How did revenue trend by region this quarter?" — analytics over billions of rowsStores columns, not rows: an aggregation over one column reads one column. Real-time inserts, SQL interface, absurd scan speed. The answer to "our analysts' queries take hours."

Notice what's not on this table: Postgres. Your relational database is still the right answer for transactional state with joins and constraints — none of the above replace it. They take the workloads Postgres does badly: writes it can't absorb, scans it can't finish, lookups it can't serve at p99.

Principle: choose the storage for the query pattern, not the hype. Every one of these systems is the best in the world at its own question and mediocre-to-terrible at the others. "Should we use Cassandra or ClickHouse?" is unanswerable until you finish the sentence: "...to do what?" Write the query first. The storage picks itself.

When to use what: five concrete scenarios

Make it mechanical. For each scenario below, the answer follows from the three shapes of work:

  1. Nightly revenue aggregates over 500 GB of click logs. Bounded input, latency budget of hours, throughput is everything → Spark batch reading Parquet, writing results back to the lake. (Post 3–5.)
  2. Fraud alerts within seconds of a suspicious payment. Unbounded input, freshness beats completeness → streaming: Kafka (you know it from the Data & Messaging track) into Structured Streaming or Flink, alert out. (Post 6.)
  3. User profile + session at login, 100k requests/sec, p99 under 10 ms. Point lookups, tiny reads → Redis in front, Cassandra/HBase as the system of record behind it. (Posts 7–9.)
  4. Analysts running ad-hoc SQL over 10 TB of events. Large scans over a few columns → ClickHouse. On big analytical scans a row-oriented system like Postgres often has to read substantially more data; a columnar engine reads only the columns the query needs and compresses similar values efficiently. (Post 9.)
  5. Exactly-once billing from a payment stream. Correctness over latency → Flink or Structured Streaming with idempotent sinks, event-time windows, watermarks for late data. (Post 6.)

If you can place a workload into one of the three boxes and name its query pattern, you've done 80% of data architecture. The remaining 20% — tuning, exactly-once semantics, schema evolution — is what posts 2 through 10 are for.

The platform map: how the pieces fit

The data platform, bottom to top SERVING Cassandra wide-column · tunable HBase wide-column · strong row Redis in-memory · ms ClickHouse columnar · OLAP PROCESS Spark batch + SQL + streaming Flink true streaming MapReduce legacy · know it PIPE Kafka the durable log between batch, streaming, and serving — Data & Messaging track RESOURCE YARN Hadoop-native scheduling Kubernetes where scheduling is moving STORAGE HDFS blocks · replication Parquet / ORC columnar files Delta Lake ACID on the lake

Read it bottom-up: storage holds the bytes (HDFS for the cluster, Parquet files for the format, Delta Lake for transactions on top). Resource managers hand out machines. The pipe — Kafka — carries events between everything. Processing engines turn bytes into answers. The serving layer answers the millisecond questions. Every post in this track lives on one row of this diagram, and the capstone wires all five rows together.

Modern reality check. HDFS teaches the foundations — and the next post goes hands-on with it — but many current data platforms store lake data in cloud object storage (S3, GCS, Azure Blob) instead of HDFS. Parquet remains the file format; table formats like Delta Lake, Apache Iceberg, or Hudi add metadata and transactional behavior above those files. The diagram's storage row reads the same either way: durable blob storage at the bottom, files in the middle, table transactions on top.

You will also encounter: Trino (a distributed SQL query engine — it queries the lake rather than owning compute the way Spark does), Hive (the original SQL-on-Hadoop layer, still underneath plenty of warehouses), Airflow / Dagster (orchestration — scheduling and dependency management for pipelines; the capstone touches this), and Iceberg / Hudi (table formats alongside Delta Lake). The diagram above is the map, not the territory — the ecosystem is wider than any one picture.

Lab: prove your environment works

Every post in this track runs real commands, so let's establish the baseline now: Spark 3.5.3 on JDK 17, local mode, one smoke test. (Why JDK 17 and not the JDK 21 from the other tracks? Spark 3.5's supported runtimes are Java 8/11/17 — on 21 it trips over module encapsulation. Data engineering means matching the tool's runtime, not the newest one.)

$ java -version
openjdk version "17.0.20.1" 2026-08-18
OpenJDK Runtime Environment Temurin-17.0.20.1+1 (build 17.0.20.1+1)

$ spark-shell --master "local[2]"
[...banner elided...]
      /___/ .__/\_,_/_/ /_/\_\   version 3.5.3
      /_/
Using Scala version 2.12.18 (OpenJDK 64-Bit Server VM, Java 17.0.20.1)

scala> spark.version
res0: String = 3.5.3

scala> val rdd = sc.parallelize(Seq(1,2,3))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:23

scala> println("SUM=" + rdd.sum())
[Stage 0:>                                                          (0 + 2) / 2]
SUM=6.0

Read that output like a Spark developer: sc.parallelize created a distributed RDD split into partitions; local[2] gave Spark two worker threads to execute tasks concurrently across those partitions. The [Stage 0: (0 + 2) / 2] line is the task tracker showing both partitions running, and SUM=6.0 is the reduced result coming back to the driver. That's the whole distributed-computing contract in four lines: split the data, run the function everywhere, combine the answers. (Partitions are worth learning early — they become central in posts 3 and 4.)

Environment note (honest, sandbox-specific). On this lab machine, Spark's SQL/DataFrame engine can't run: its internal network transport hits a sandbox networking limitation and dies with a Netty "too large frame" error. Plain RDD operations run fine, which is all this smoke test needs. On any normal machine the full Spark SQL session works; the DataFrame-heavy labs in posts 3–6 were designed against that, and each states its environment plainly.

Where this track goes

The map is drawn; now we walk it. Each post is one row of the platform diagram, with real labs:

  1. This post — the landscape: batch vs streaming vs serving, and why the JVM runs it all.
  2. HDFS & MapReduce: The Foundations — blocks, replication, YARN containers, and a real MapReduce job in Java. The history you need to understand everything else.
  3. Spark Core: RDDs, DataFrames & the Java API — lineage, transformations vs actions, and the Java-specific sharp edges (lambdas, Encoders, serialization).
  4. Spark SQL Deep Dive — Catalyst, predicate pushdown, partitioning, and reading EXPLAIN output like a query plan.
  5. File Formats & the Lakehouse — Parquet internals, partitioning schemes, Delta Lake ACID, schema evolution.
  6. Streaming: Spark Structured Streaming — event time, watermarks, windows, and exactly-once sinks.
  7. Cassandra: Wide-Column at Scale — data modeling query-first, tunable consistency, CQL in Java.
  8. HBase: Strong Consistency on Hadoop — regions, row-key design, and when strong consistency is worth the coordination cost.
  9. Redis & ClickHouse: Storage & Serving — sub-millisecond lookups and columnar analytics, two very different answers to downstream access.
  10. Capstone: End-to-End Pipeline — all five rows wired together: ingest → process → serve, orchestrated and observable.

Cheat sheet

  • OLTP vs OLAP is the original split: one database asked to do both badly is what creates the need for a data platform.
  • Batch asks "what happened?", streaming asks "what's happening?", serving asks "what about this one?" — one question per workload.
  • HDFS stores (blocks, 3x replication), YARN schedules (containers), MapReduce was the first engine (map → shuffle → reduce, disk-bound).
  • Spark won by avoiding MapReduce's per-stage disk materialization — pipelined DAG execution, cached reuse for iterative work, spilling when needed — plus one unified engine; the DataFrame API lets Catalyst optimize your queries.
  • Flink is the purpose-built streaming engine: event time, watermarks, exactly-once state. Rule of thumb, not a law: seconds-to-minutes latency → Structured Streaming is often enough; consistently low latency with deep stateful streaming → Flink.
  • The JVM runs data infra through a combination of portable runtime, mature concurrency, the JIT, tuned GCs, the library network effect, and twenty years of production debuggability.
  • Cassandra (tunable consistency, write-heavy), HBase (strong row consistency, on HDFS), Redis (in-memory speed layer), ClickHouse (columnar OLAP) — four storage/serving answers, four different access patterns.
  • Choose the storage for the query pattern, not the hype. Write the query first; the storage picks itself.
  • Modern lakes usually sit on object storage (S3/GCS/Blob), not HDFS — HDFS teaches the foundations; Parquet + a table format (Delta/Iceberg/Hudi) sit above either.
  • Kafka is the pipe between the rows of the platform diagram — you already know it from the Data & Messaging track.
  • Spark 3.5 wants JDK 8/11/17 — match the tool's runtime, not the newest JDK.

What's next

You can now place any data tool on the map and name the question it answers. But the map has a blank spot at its foundation: you haven't touched HDFS or written a MapReduce job, and every abstraction above them leaks that foundation sooner or later — Spark's shuffle is MapReduce's shuffle, and partitioning schemes only make sense once you've felt a hot key.

The next post, HDFS & MapReduce: The Foundations, goes hands-on with the old growth forest: block placement and replication you can observe, a YARN container you can watch, and a real Java MapReduce job — run end to end, slow shuffle and all, so you understand exactly what Spark was built to escape.

Field check before you move on: (1) Run the smoke test above on your own machine with local[4] instead of local[2] — watch the stage tracker line and explain what changed. (2) Pick one system from the NoSQL table and write down, in one sentence, the query pattern it's built for and one query pattern it would serve badly. (3) Take the five scenarios from "When to use what" and argue for a different answer for one of them — under what changed assumption would your alternative win? (4) Find one workload in your current job or a past project and place it in batch/streaming/serving — then check whether the tool it actually uses matches the box.

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