Spark Core: RDDs, DataFrames & the Java API
Post 1 drew the map and gave you a four-line Spark smoke test. Now it's time to earn those four lines. That nightly revenue report — the one that takes four hours and times out — is a batch job: read a few hundred gigabytes of orders, filter, group, sum, write. This post is the machinery underneath it: the RDD, the unit Spark thinks in; laziness, the reason your program does nothing until the last line; partitions, the unit it scales in; and caching, the one performance lever that is entirely your responsibility. Every claim with a transcript below ran in this track's lab — Spark 3.5.3, local[2], Temurin JDK 17 — and the two sections that couldn't run here are labeled honestly instead of faked.
Lab honesty, up front. Everything in the RDD sections ran for real in spark-shell --master "local[2]" on this track's lab machine: the toDebugString lineage, the laziness demonstration, the partition counts with their stage-tracker lines, and the caching experiment with accumulator values — all pasted verbatim. Two sections could not run here and are marked [NOT RUN]: the Java-API code (lambdas, serialization, Encoders) needs a compiled jar shipped by spark-submit, and the DataFrame examples need the SQL engine — both hit this sandbox's Netty transport failure. Those sections show complete, careful code and describe standard documented behavior; nothing in them is presented as observed output.
The abstraction: an RDD
An RDD — Resilient Distributed Dataset — is Spark's core abstraction: an immutable, partitioned collection of records spread across the cluster, with a recipe for rebuilding any lost piece. Read that sentence again, because each word carries a design decision:
- Immutable. You never modify an RDD; every operation returns a new one. Immutability is what makes the next two words possible.
- Partitioned. The dataset is chopped into partitions — the unit of parallelism. A task processes one partition. More partitions (up to a point) means more parallelism; fewer means less scheduling overhead. This is the single most load-bearing concept in Spark, and this post returns to it in every section.
- Resilient — through lineage, not replication. HDFS replicates every block 3x so a dead machine loses nothing. Spark doesn't replicate data at all. Instead, every RDD remembers the recipe that built it: "partition 3 of
big= filter of partition 3 ofdoubled= map of partition 3 ofnums." If an executor dies, Spark replays the recipe for just the lost partitions. That recipe is called lineage, and you can read it directly.
Lineage you can read
Build a small chain and ask Spark what it remembers:
scala> val nums = sc.parallelize(1 to 20, 4)
nums: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:23
scala> val doubled = nums.map(_ * 2)
doubled: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[1] at map at <console>:23
scala> val big = doubled.filter(_ > 10)
big: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[2] at filter at <console>:23
scala> println(big.toDebugString)
(4) MapPartitionsRDD[2] at filter at <console>:23 []
| MapPartitionsRDD[1] at map at <console>:23 []
| ParallelCollectionRDD[0] at parallelize at <console>:23 []
scala> println("PARTITIONS=" + big.getNumPartitions)
PARTITIONS=4
That indented tree is the lineage: big is a filter of a map of a parallelized collection, and the leading (4) says it currently has 4 partitions. No data has moved anywhere yet — and that is the next section's whole point. Notice also how the partition count survived two transformations: map and filter are one-to-one with their input, so they inherit its partitioning. Transformations that reshuffle data across partitions (like repartition) change it; you'll see that in the stage tracker shortly.
Transformations vs actions: nothing runs until you ask
Look at the transcript above once more: three lines built three RDDs, and not one stage ran. No [Stage 0: ...] tracker lines, no tasks, no data movement. Spark recorded the recipe and waited. Transformations (map, filter, flatMap, repartition...) build the lineage DAG. Actions (count, collect, sum, saveAsTextFile...) execute it. This laziness lets Spark plan the whole pipeline — and, for DataFrames, optimize it — before touching data.
The sharpest way to feel it is to build a chain that cannot succeed:
scala> val lazyChain = sc.parallelize(1 to 5).map(x => 1 / (x - 3))
lazyChain: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[4] at map at <console>:23
scala> println("chain built; if lazy, no exception yet")
chain built; if lazy, no exception yet
scala> lazyChain.collect()
[Stage 0:> (0 + 2) / 2]
org.apache.spark.SparkException: Job aborted due to stage failure: Task 1 in stage 0.0
failed 1 times, most recent failure: Lost task 1.0 in stage 0.0 (TID 1)
(198.19.0.2 executor driver): java.lang.ArithmeticException: / by zero
[...stack trace elided...]
Two lessons in one transcript. First: the divide-by-zero lived in the map line, but nothing complained until collect() — the action — forced execution. Laziness defers not just work but failure, which is why Spark stack traces point at the action line even when the bug is three transformations up. Second: the executor's ArithmeticException arrived wrapped in Spark's own SparkException: Job aborted due to stage failure. Read these inside-out: the cause at the end is your bug; the wrapper tells you which stage and task died. (My transcript also caught me mid-demo: a narrow catch for ArithmeticException missed, because what the driver actually throws is the SparkException wrapper — laziness means the failure crosses a process boundary before it reaches you.)
| Transformations (lazy — build the recipe) | Actions (eager — run the recipe) |
|---|---|
map, flatMap, filter | collect (bring all data to the driver) |
repartition, coalesce | count, sum, reduce |
union, distinct, sortBy | take(n), first (bring a little) |
groupByKey, reduceByKey, join | foreach, saveAsTextFile (side effects) |
A rule of thumb worth internalizing: if the return type is still an RDD, it was lazy; if it returns anything else, it ran. It holds for the whole core API, and it will save you from the classic beginner bug — calling collect() "just to check" on a terabyte dataset and watching the driver run out of memory.
Partitions: the unit of parallelism
Partitions deserve their own section because nearly every Spark performance story is a partition story. The rule of thumb: one task processes one partition, so your partition count is your parallelism ceiling. Two hundred cores and ten partitions means one hundred and ninety cores idle. Ten thousand tiny partitions means the scheduler spends more time scheduling than computing. You will tune this number on every real job.
Where do partition counts come from? parallelize takes one explicitly (we used 4 above); file reads default to roughly one partition per input split; and two transformations change the count — repartition and coalesce. They look similar and behave very differently, and the stage tracker shows it:
scala> println("BEFORE=" + big.getNumPartitions)
BEFORE=4
scala> val rp = big.repartition(8)
rp: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[8] at repartition at <console>:23
scala> println("AFTER-REPARTITION-PARTITIONS=" + rp.getNumPartitions)
AFTER-REPARTITION-PARTITIONS=8
scala> println("REPART-COUNT=" + rp.count())
[Stage 1:==============> (1 + 2) / 4]
[Stage 2:> (0 + 2) / 8]
REPART-COUNT=15
scala> val co = rp.coalesce(2)
co: org.apache.spark.rdd.RDD[Int] = CoalescedRDD[9] at coalesce at <console>:23
scala> println("AFTER-COALESCE-PARTITIONS=" + co.getNumPartitions)
AFTER-COALESCE-PARTITIONS=2
scala> println("COALESCE-COUNT=" + co.count())
COALESCE-COUNT=15
Read it like a Spark developer. repartition(8) is a wide dependency: every output partition needs data from every input partition, so Spark inserts a shuffle — visible as two stages, Stage 1 reading the original 4 partitions and Stage 2 writing the new 8. A shuffle moves data across the network, and on a real cluster it is usually the most expensive thing in your job. coalesce(2), going down in partitions, is a narrow dependency (CoalescedRDD, not a shuffle): each output partition is assembled from a few local input partitions, no cross-network redistribution. The rule of thumb, not a law: repartition when you need more parallelism or an even spread (it shuffles, and that's the price); coalesce when you're done and want fewer, larger output files (it avoids the shuffle). And coalesce with shuffle=true exists for the case where avoiding the shuffle gives you badly skewed partitions — the exception that proves the rule is a choice, not a default.
Caching: explicit, or it didn't happen
Here's the misconception to kill early: Spark does not keep your data in memory. Left alone, it keeps nothing — every action replays the lineage from scratch, pipelining transformations stage by stage and spilling to disk when a stage's working set outgrows memory. That recompute is the price of resilience-by-lineage, and most of the time it's the right trade. But when you reuse an expensive RDD — an iterative algorithm's dataset, a filtered fact table read by five downstream queries — replaying the recipe five times is pure waste. The fix is explicit: persist() (or its alias cache()).
An accumulator makes the invisible visible. It counts how many times the map function actually executes across actions:
scala> val acc = sc.longAccumulator("mul-count")
scala> val raw = sc.parallelize(1 to 100, 4).map(x => { acc.add(1); x * 2 })
scala> raw.count()
res11: Long = 100
scala> println("after action1 acc=" + acc.value)
after action1 acc=100
scala> raw.count()
res13: Long = 100
scala> println("after action2 (no cache) acc=" + acc.value)
after action2 (no cache) acc=200
scala> raw.persist()
scala> raw.count()
[Stage 7:> (0 + 2) / 4]
res16: Long = 100
scala> println("after action3 (persisted) acc=" + acc.value)
after action3 (persisted) acc=300
scala> raw.count()
res18: Long = 100
scala> println("after action4 (from cache) acc=" + acc.value)
after action4 (from cache) acc=300
scala> println("storage=" + raw.getStorageLevel)
storage=StorageLevel(memory, deserialized, 1 replicas)
The accumulator is the whole lesson in four numbers. Actions 1 and 2, no cache: 100 → 200 — Spark recomputed the entire map for the second action without telling you. Action 3, after persist(): 200 → 300 — the data was computed once more and then written into the cache (persist marks the RDD; the next action materializes it). Action 4: 300 → 300 — zero recomputation, served from the cache. Note the storage level the lab reported: StorageLevel(memory, deserialized, 1 replicas) — the RDD default is MEMORY_ONLY: deserialized Java objects in executor memory. (The DataFrame/Dataset API's documented default is MEMORY_AND_DISK; the two APIs differ here, which is worth knowing when you move between them in post 4.)
Principle: if you didn't persist it, Spark recomputed it. Caching is opt-in, never automatic. And a cache is not a guarantee of residency: with MEMORY_ONLY, partitions that don't fit are simply dropped and recomputed on demand — not spilled. Spilling to disk is what MEMORY_AND_DISK (and the shuffle path, which always spills) do when memory runs out. So the honest mental model has three layers: pipelined execution (the default — stream through stages, keep almost nothing), explicit caching (you name what survives, with a storage level that says what "survives" means), and spilling (the safety valve when memory runs out). "Spark is in-memory" collapses all three into a slogan; the lab shows you which one is actually the default.
The Java API's sharp edges [NOT RUN]
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
So far the labs have been Scala, because spark-shell is a Scala REPL. In production you'll write the Java API, and it has three sharp edges that bite every Java developer exactly once. The code below is complete and careful — it just couldn't execute in this sandbox, which can't ship a compiled jar over its broken transport.
Lambdas and serialization: the NotSerializableException
Every lambda you pass to a Spark transformation is shipped to the executors — which means it, and everything it captures, must be serializable. The classic failure looks innocent:
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
public class OrderLabelJob {
// java.text.SimpleDateFormat is NOT serializable
private final SimpleDateFormat fmt = new SimpleDateFormat("yyyy-MM-dd");
// assume Order.getTimestamp() returns java.util.Date
public JavaRDD<String> labelDates(JavaSparkContext sc, JavaRDD<Order> orders) {
// the lambda captures `this` (to reach fmt) -> `this` is not serializable
return orders.map(o -> fmt.format(o.getTimestamp()));
}
}
At the action, the driver tries to serialize the closure, walks its captured references, finds the enclosing OrderLabelJob instance (captured so the lambda can reach fmt), and fails: the documented symptom is a Task not serializable job failure caused by java.io.NotSerializableException. The rule of thumb: a lambda may only capture serializable state — and an instance-method lambda always captures this. Three standard fixes, in increasing order of preference:
- Create the unserializable object inside the lambda (per record — simple, but pays construction cost per element).
- Use
mapPartitions— construct it once per partition, on the executor, where it never crosses the wire:JavaRDD<String> labels = orders.mapPartitions(it -> { SimpleDateFormat f = new SimpleDateFormat("yyyy-MM-dd"); // built on the executor List<String> out = new ArrayList<>(); while (it.hasNext()) { out.add(f.format(it.next().getTimestamp())); } return out.iterator(); }); - Make the holder serializable or static — a
static finalfield isn't part of the captured instance, so it never gets serialized with the closure. (AThreadLocalstatic is the usual thread-safety answer for formatters.)
Debugging tip from the trenches: when the exception names a class you didn't expect — often your outer job class or an anonymous inner class — that's the captured this. Make the lambda static-friendly or move it out of the instance.
Encoders and the JavaBean contract
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
RDDs of JVM objects serialize with Java/Kryo serialization — opaque blobs Spark can't inspect. Datasets need better: to optimize, Spark must understand your object's fields. The bridge is an Encoder — a serializer that also exposes the schema — and for Java beans you get one via Encoders.bean(), provided your class honors the JavaBean contract:
public class Order implements Serializable {
private String orderId;
private double amount;
public Order() {} // 1. public no-arg constructor
public String getOrderId() { return orderId; } // 2. getters...
public void setOrderId(String orderId) { this.orderId = orderId; } // ...and setters
public double getAmount() { return amount; }
public void setAmount(double amount) { this.amount = amount; }
// 3. (implied) field types Spark can map: String, double, long, java.sql.Timestamp, ...
}
Encoder<Order> orderEncoder = Encoders.bean(Order.class);
Dataset<Order> orders = spark.read()
.parquet("s3://lake/orders/")
.as(orderEncoder); // Row -> Order, with schema
Dataset<Order> bigOrders =
orders.filter((FilterFunction<Order>) o -> o.getAmount() > 1000);
The contract is strict because the encoder is generated from it by reflection: miss the no-arg constructor and encoder creation fails; a property without both a getter and a setter won't survive the round trip — when a field comes back null or default after a Dataset operation, the bean contract is the first suspect. The rule of thumb: if Encoders.bean() can't introspect it, Catalyst can't optimize it — which is the doorway to the DataFrame API.
DataFrames: the same engine, a smarter handle [NOT RUN]
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
Everything above was the RDD API: you hand Spark opaque JVM objects and hand-written functions, and it executes exactly what you wrote. A DataFrame (Dataset<Row> in the typed API) changes the deal: instead of objects, you work with named, typed columns — and because Spark can see the structure of your computation, it can rewrite it:
Dataset<Row> df = spark.read()
.option("header", "true")
.csv("s3://lake/orders.csv");
df.select("order_id", "amount", "region")
.filter("amount > 1000")
.groupBy("region")
.sum("amount")
.show();
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
Conceptually, your code passes through the Catalyst optimizer in four phases: the query is parsed into a logical plan, analyzed against the catalog (do these columns exist?), optimized with rewrite rules, then turned into a physical plan of RDD operations. The two rewrites to know by name are predicate pushdown (a filter is moved as close to the data source as possible, so irrelevant rows are never read) and column pruning (only the selected columns are read — the reason columnar files like Parquet pair so well with DataFrames). Nothing here changes the execution model you learned above — a DataFrame job still runs as stages over partitions, still shuffles on wide dependencies, still recomputes unless you persist. What changes is that Spark gets a vote on how your logic executes.
Rule of thumb, not a law: reach for DataFrames when your logic is expressible as rows, columns, and relational operations — which, for analytics, is most of the time. Reach for RDDs when you need per-partition control (mapPartitions), custom partitioning, or truly opaque objects. The execution details — reading EXPLAIN output, partitioning for joins, and why one plan beats another — are post 4's whole subject.
Cheat sheet
- An RDD is an immutable, partitioned collection with lineage: resilience through the recipe, not replication.
toDebugStringprints the lineage — read it when a job misbehaves; it's the first thing a Spark veteran asks for.- Transformations are lazy (return RDDs); actions are eager (return anything else). Nothing — including failures — happens until an action.
- Executor exceptions arrive wrapped in
SparkException: Job aborted due to stage failure; the cause at the end is your bug. - Partitions are the unit of parallelism: one task per partition. Tune the count; it's the most common performance lever.
repartition= wide dependency = shuffle = two stages (expensive, and the price of more/even partitions).coalescedown = narrow = no shuffle.- If you didn't persist it, Spark recomputed it. Caching is explicit and opt-in; the RDD default is MEMORY_ONLY.
- Three-layer model: pipelined execution (default) → explicit caching (your choice, your storage level) → spilling (the safety valve). "In-memory" is a slogan, not the default.
- Java lambdas ship to executors: capture only serializable state; instance-method lambdas capture
this.mapPartitionsbuilds unserializable helpers on the executor. Encoders.bean()needs the full JavaBean contract: public no-arg constructor, getters, setters, mappable field types — a property missing either accessor won't survive the round trip.- DataFrames give Catalyst a structured plan to optimize (predicate pushdown, column pruning); the execution model underneath is still stages over partitions.
What's next
You now own Spark's core: lineage you can read, laziness you can feel, partitions you can count, and caching you control. But the DataFrame section ended on a promise — Catalyst rewrites your query, and you haven't yet seen how it decides, or how to tell when its choice is bad. That's a skill, not trivia: on real jobs the difference between a plan that prunes and a plan that scans is the difference between minutes and hours.
The next post, Spark SQL Deep Dive, opens the optimizer's black box: reading EXPLAIN output like a query plan, partitioning and bucketing for joins, and the physical operators your DataFrame code actually becomes. Bring the partition intuition from this post — it's the foundation everything in post 4 stands on.
Field check before you move on: (1) Rebuild the lineage lab with your own chain of three transformations, print toDebugString, and sketch the DAG on paper — then predict the stage count before you run an action, and check. (2) Rerun the accumulator experiment but call persist() before the first action; predict the four accumulator values, then compare with what you saw here (100/200/300/300) and explain the difference. (3) repartition a 4-partition RDD to 16 and run count(): how many stages appear in the tracker, and why two and not one? (4) Take any lambda you've written for a Spark transformation and list everything it captures — the enclosing instance, every field it touches — then mark each serializable or not.
Continue: Java Learning Roadmap 2026
Comments
Post a Comment