In-Memory Model and RDDs

Spark's In-Memory Model and the RDD Abstraction

Spark 129 began in 2009 at UC Berkeley's AMPLab under Matei Zaharia and became a top-level Apache project in 2014. Its "Resilient Distributed Datasets" paper (NSDI 2012) showed that keeping data in memory speeds iterative work by an order of magnitude; in 2014 Spark sorted 100 TB in 23 minutes on 206 machines, against Hadoop 129 's 72 minutes on 2,100.

An RDD is an immutable collection split into partitions across the cluster. Transformations (map, filter, reduceByKey) only describe a new RDD; an action (count, collect) runs the work. Each RDD records its lineage, so a lost partition is rebuilt by replaying the chain for that partition alone. The listing runs in the PySpark 129 shell (spark-shell and pyspark), where spark exists, over BookNest's order history: JSON, Columnar and Binary Formats's seeded generator run at ten times its default size, one million sample orders in data/raw/orders.jsonl.

An RDD pipeline with lineage and an in-memory cacheJavaScript
import json, time
from operator import add
sc = spark.sparkContext
lines = sc.textFile("data/raw/orders.jsonl", 4)            # at least 4 partitions
done = lines.map(json.loads).filter(lambda o: o["status"] == "delivered").cache()
revenue = (done.flatMap(lambda o: o["items"])               # lazy: nothing runs yet
           .map(lambda i: (i["book_id"], i["qty"] * i["unit_price"])).reduceByKey(add))
print(revenue.toDebugString().decode())
for run in (1, 2):                                          # run 2 reads the cache
    t0 = time.perf_counter()
    n = done.count()
    print(f"count {run}: {n:,} orders in {time.perf_counter() - t0:.2f}s")
Output
(7) PythonRDD[7] at RDD at PythonRDD.scala:59 []
...
 +-(7) PairwiseRDD[4] at reduceByKey at /home/dev/v7-l3/ch05/listings/run.py:13 []
...
    |  data/raw/orders.jsonl HadoopRDD[0] at textFile at DirectMethodHandleAccessor.java:103 []
count 1: 890,703 orders in 18.81s
count 2: 890,703 orders in 2.92s

Read the lineage bottom up: +- marks the shuffle that reduceByKey forces, and (7) is the partition count (the 238 MB file splits into seven blocks of about 32 MB). The second count() reads cached partitions and runs five to six times faster.