Apache Spark — The answer to a slow job is in the plan and the event log
A DataFrame is a recipe, not a result, so it is cooked again every time
In one line
A DataFrame is not a computed result but a lineage from the source to the result, so every time you call an action, it recomputes from the source. A cache holds on to the result of the first computation, and a checkpoint writes the result to files and cuts the lineage.
Why is the same computation done twice
Suppose you call count() on one DataFrame that reads a CSV, cleans it and joins it, and then write the same thing to files. Looking only at the code, the cleaning happens once. In reality it happens twice. In the event log, both jobs start from the CSV scan.
The reason is that transformations are lazy. filter or join only adds one line to the plan and computes nothing. When an action is called, Spark runs that plan from the source to the end, hands the result to the driver or to files, and then throws away the intermediate results. The next action starts again from the source. This is not a defect but a design. If it held on to all intermediate results, memory would fill up quickly, and with the lineage alone a lost piece can be rebuilt at any time.
So choosing only the intermediate results that are used several times and holding on to them is up to people. That is the cache.
How it works
cache() is also lazy. At the moment you call it, only a mark is left saying "hold on to this result when it is first computed". When the first action runs, each partition is stored as it is computed, and subsequent actions start from InMemoryTableScan in the plan instead of the source scan. As the RDD programming guide explains, it is stored when first computed, and subsequent actions reuse it.
from pyspark import StorageLevel
clean = raw.filter("qty > 0").join(products, "product_id")
clean.cache() # 표시만 한다 — 아직 아무것도 저장되지 않았다
clean.count() # 첫 행동: 계산하면서 저장
clean.write.parquet("/tmp/out") # 계획이 InMemoryTableScan 에서 시작한다
clean.unpersist() # 다음 행동은 다시 원본부터
slim = raw.select("order_id", "qty").persist(StorageLevel.DISK_ONLY)
The storage level decides where to put it. The default is the easy place to get confused. The cache() of an RDD is MEMORY_ONLY, which keeps it in memory only. In contrast, the cache() of a DataFrame is MEMORY_AND_DISK_DESER, and partitions that do not fit in memory are kept on disk. The documentation says this default was changed in 3.0 to match Scala. If you want a different level, you pass it to persist(), but you cannot give a new level to a DataFrame whose storage level is already set. You have to unpersist() first.
A DataFrame cache is placed in memory not row by row but in columnar format. The performance tuning documentation explains that it reads only the needed columns from a cached table and tunes compression by itself.
On which level to choose, the RDD guide's advice is short. If it fits comfortably in memory, leave the default, and do not spill to disk unless the computation is expensive or filters out a large amount. This is because recomputing can be as fast as reading from disk. If the source is a single local CSV and the transformations are light, a DISK_ONLY cache gains almost nothing.
Where the cache lives: unified memory
A cache is not placed in free space. According to the tuning guide, the execution memory used by shuffles, joins, sorts and aggregations and the storage memory used by the cache share one region. The size of that region is spark.memory.fraction (default 0.6), a ratio of the JVM heap minus 300MiB. Within it, the share spark.memory.storageFraction (default 0.5) is the part that execution memory cannot push the cache out of.
Estimating with this lab's driver memory of 1g gives (1024 - 300) × 0.6, that is, a bit over 400MiB shared by execution and the cache. If you cache a large table, the space for joins and sorts shrinks by that much, and if the cache overflows, then as the RDD guide says, the least recently used partitions (LRU) are evicted first. An evicted partition is recomputed along the lineage the next time it is needed. So a cache never produces a wrong answer. It just becomes slower.
unpersist and checkpoints
A cache you are done with is released with unpersist(). The RDD guide says that this call by default does not wait. To wait until the resources are actually released, you give the blocking argument. An action after releasing it starts again from the source scan in the plan. One more thing: a DataFrame cache is shared by all sessions of the cluster. If another session has cached the same plan, my query uses it too.
A cache leaves the lineage as it is. That is also why a lost partition can be rebuilt. But when you pile dozens of transformations onto the same DataFrame, as in an iterative algorithm, the lineage itself becomes a problem. The checkpoint documentation explains that in such cases the plan can grow exponentially and that a checkpoint cuts off the logical plan. It writes the result as files to the directory set by SparkContext.setCheckpointDir() or spark.checkpoint.dir, and the plan after it starts not from the source but from that file. In the plan, instead of the source scan, a node that reads an existing RDD appears. The default is eager execution, so a job runs at the moment you call it.
localCheckpoint() stores to the executors' caching system instead of files. It is fast, but as the documentation itself says, it is not reliable: if you lose an executor, the lineage has also been cut, so there is no way to rebuild it.
SQL has the same tool. Unlike the cache() of a DataFrame, CACHE TABLE is eager by default, and you have to add LAZY for it to cache on first use. If you do not give a storage level, it is MEMORY_AND_DISK.
What it looks like in the field
First, caching something used only once. If you cache a DataFrame that has only one action, you just pay the extra cost of storing. Put a cache only at the forks that are used two or more times.
Second, holding on to a cache and forgetting it. When caches pile up inside a long notebook or a service, execution memory shrinks and the shuffle starts spilling to disk. Release it when you are done.
Third, you thought you cached it but it is not in the plan. If you call an action on a different DataFrame that has a new transformation attached after cache(), only the cached part inside that DataFrame's plan is reused. To check that the cache was really used, find InMemoryTableScan in explain() after the action. The Storage tab of the web UI also shows the cache only after it has materialized.
What really matters in practice
- A DataFrame is a recipe. Every action recomputes from the source.
- cache() is lazy. It is stored when the first action runs.
- The default storage level is MEMORY_AND_DISK_DESER for a DataFrame and MEMORY_ONLY for an RDD.
- A cache shares one region with execution memory. When you are done, unpersist.
- If the lineage is too long, cut it with a checkpoint. A cache does not cut the lineage.
What you will do in the next lab
Without a cache, you call two actions and confirm from the event log and the plan that both jobs read the CSV again. After cache(), you see that the plan of the second action starts from InMemoryTableScan, look at the storage level, and use persist to specify a level that keeps it only on disk. You confirm that whether it is cached changes before and after unpersist, cut a DataFrame with ten transformations piled on by a checkpoint and see whether the same answer comes out, and then do the same thing with SQL's CACHE TABLE. Finally, in the event log of an application with block update logging on, you add up the sizes of the cache blocks and measure in numbers how much memory the cache took.