TT Lab
Get started
Learn Learning paths Courses

Apache Spark — The answer to a slow job is in the plan and the event log

Avoid computing the same thing twice — how a cache shows up in the plan

Continue in TT Lab

Goal

When you call two actions on the same DataFrame, you confirm from the plan that without a cache it recomputes from the source, and see how cache, persist, unpersist, checkpoints and CACHE TABLE appear in the plan and the event log. You also measure with block update events how much memory the cache actually eats.

Why it matters

A DataFrame is not a result but a way of making it (the lineage). Every time you call an action, Spark runs that method again from the beginning. Code that uses the same intermediate result twice reads the source twice and does the join twice, without an error or a warning. cache() is a mark telling Spark to keep that intermediate result on the executors. It is only a mark, so the first action fills it, and subsequent actions read the cache with InMemoryTableScan in the plan instead of reading the source. The storage level (persist) chooses memory, disk and serialization, and decides what to give up when memory runs short. A cache you are done with must be put down with unpersist so that it becomes memory for other work. If the lineage gets too long (iterative computation), the plan itself grows and the driver slows down. A checkpoint writes the result to files and cuts the lineage, making the plan after it a single line: "read this file".

Steps

  1. Put a function that attaches an amount to the paid orders and a function that calls two actions (count, sum(amount)) in /root/spk/cache/common.py, and with /root/spk/cache/none.py (application spk-cache-none), write the results of the two actions without a cache to /root/spk/cache/out/none.json.
  2. In /root/spk/cache/cached.py (application spk-cache-on), apply cache() to the same DataFrame and write the results of the two actions to /root/spk/cache/out/cached.json and the storage level (str(df.storageLevel)) to /root/spk/cache/out/cached_level.txt.
  3. In /root/spk/cache/disk.py (application spk-cache-disk), do the same with persist(StorageLevel.DISK_ONLY) and write the result to /root/spk/cache/out/disk.json and the storage level to /root/spk/cache/out/level.txt.
  4. In /root/spk/cache/unpersist.py (application spk-cache-un), run it in the order cache → action → unpersist() → action, and write the before and after of whether it is cached (df.is_cached) to /root/spk/cache/out/unpersist.json as {"before": 참거짓, "after": 참거짓} (both values are booleans).
  5. In /root/spk/cache/ckpt.py (application spk-cache-ckpt), set the checkpoint folder to /root/spk/cache/ckpt, and with the result of checkpoint() after ten withColumn calls that add 1 to 10 to the amount in turn, write the sum of the amount per channel as Parquet (channel,amount) to /root/spk/cache/out/ckpt_result.
  6. In /root/spk/cache/cache_sql.py (application spk-cache-sql), create the temporary view paid, run CACHE TABLE paid_cached AS SELECT channel, amount FROM paid, compute the sum per channel, and write to /root/spk/cache/out/cache_table.json as {"is_cached": 참거짓, "by_channel": {채널: 합}} (is_cached is a boolean, and by_channel maps each channel name to its sum).
  7. With /root/spk/cache/size.py (application spk-cache-size, spark.eventLog.logBlockUpdates.enabled=true), fill the cache, then add up the rdd_ block sizes in that log and write them to /root/spk/cache/out/cache_size.json as {"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수} (all three values are integers).
  8. In /root/spk/cache/report.md, write three sections: ## 다시 계산하는 비용, ## 저장 수준 and ## 캐시의 크기 (use exactly these Korean headings in this order; they mean "Cost of recomputing", "Storage level" and "Size of the cache"). In the third section, put the number of blocks and the memory bytes from step 7.

Notes

Two actions without a cache: the source is read twice

In /root/spk/cache/common.py, put a function that joins the paid orders with the products and attaches amount (qty×price), and a function that calls the two actions count() and sum(amount) on a DataFrame and writes {"count": 정수, "sum": 정수} (two integers) to a file. Create /root/spk/cache/none.py with the application name spk-cache-none and write the result without a cache to /root/spk/cache/out/none.json.

Two actions mean two SQL executions, and both have Scan csv at the very bottom of the plan. This means the same file is read twice and the same join is done twice. The grader checks whether there are two or more executions that read the source without a cache.

cache(): the first action fills it and the next one reads it

Create /root/spk/cache/cached.py with the application name spk-cache-on, apply cache() to the step 1 DataFrame, and write the results of the same two actions to /root/spk/cache/out/cached.json and str(df.storageLevel) to /root/spk/cache/out/cached_level.txt.

InMemoryTableScan appears in the plan of the second execution. Also notice that the default storage level of a DataFrame differs from the cache of an RDD: if memory runs short, it spills to disk. The result numbers must be the same as in step 1.

persist(DISK_ONLY): choosing the storage level

Create /root/spk/cache/disk.py with the application name spk-cache-disk, keep the step 1 DataFrame with persist(StorageLevel.DISK_ONLY), and write the results of the same two actions to /root/spk/cache/out/disk.json and str(df.storageLevel) to /root/spk/cache/out/level.txt.

Keeping it only on disk saves executor memory, but it needs deserialization every time you read it. See that the StorageLevel printed on the InMemoryRelation line of the plan has changed. Even with the name InMemory, it can be on disk.

unpersist: if you put the cache down, it is recomputed

Create /root/spk/cache/unpersist.py with the application name spk-cache-un, run the step 1 DataFrame in the order cache() → count() → unpersist() → sum(amount), and write df.is_cached before and after unpersist() to /root/spk/cache/out/unpersist.json as {"before": 참거짓, "after": 참거짓} (both values are booleans).

A cache is not free. It takes executor memory and reduces the space for other work. Put down a cache you are done with. The grader checks from the log whether an execution that does not read the cache follows the execution that read the cache.

Cut the lineage with a checkpoint

Create /root/spk/cache/ckpt.py with the application name spk-cache-ckpt, set the checkpoint folder to /root/spk/cache/ckpt, and with the result of checkpoint() after ten withColumn calls that add 1, 2, …, 10 in turn to amount of the step 1 DataFrame, write the sum of amount per channel as Parquet (channel,amount) to /root/spk/cache/out/ckpt_result.

checkpoint() writes the result to the checkpoint folder on the spot, and the plan of the DataFrame it returns becomes a single line that reads that file (Scan ExistingRDD) rather than the source. A cache keeps the lineage (if lost, it is recomputed), and a checkpoint cuts the lineage (the file is the source).

SQL's CACHE TABLE is not lazy

Create /root/spk/cache/cache_sql.py with the application name spk-cache-sql, register the step 1 DataFrame as the temporary view paid, run CACHE TABLE paid_cached AS SELECT channel, amount FROM paid, then compute the sum of amount per channel from paid_cached and write to /root/spk/cache/out/cache_table.json as {"is_cached": spark.catalog.isCached("paid_cached"), "by_channel": {채널: 합}} (by_channel maps each channel name to its sum).

The cache() of a DataFrame only marks, but CACHE TABLE fills it on the spot (with CACHE LAZY TABLE, it defers). So the CACHE TABLE statement itself launches a job. The plan that reads a named cache prints Scan In-memory table paid_cached. The grader checks whether such an execution exists and the sums.

Measure the memory the cache actually eats

Create /root/spk/cache/size.py with the application name spk-cache-size and the setting spark.eventLog.logBlockUpdates.enabled=true, apply cache() to the step 1 DataFrame and fill it with count(). From that application's log, among the SparkListenerBlockUpdated events, keep only the last one per block for those whose block ID starts with rdd_, and write the count and the sums of Memory Size and Disk Size to /root/spk/cache/out/cache_size.json as {"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수} (all three values are integers).

There is one cache block per partition (rdd_<RDD 번호>_<파티션>, where the placeholders are the RDD number and the partition). A cache unpacked in memory (deserialized) can be larger or smaller than the original CSV, depending on the chosen columns and types. The habit of looking at this number before growing a cache protects executor memory.

Where to put a cache and where to put it down

In /root/spk/cache/report.md, write three sections: ## 다시 계산하는 비용, ## 저장 수준 and ## 캐시의 크기 (use exactly these Korean headings in this order; they mean "Cost of recomputing", "Storage level" and "Size of the cache"). Put the storage level strings from steps 2 and 3 in the second section and the number of blocks and the memory bytes from step 7 in the third section.

In the first section, write what was done twice when there was no cache; in the second, what the two storage levels trade off; and in the third, what percent of the executor heap (1g) that memory is.