skip to content

When several downstream Spark jobs read one expensive DataFrame, how do you choose between persist() and writing it out?

level: principalimportance: should knowfreq 38%

answer

  1. ask who the readers are first
  2. one lives inside the application
  3. the other outlives it
  4. memory is tactical, storage is contractual

basics

~20 s

Scope decides it. persist() reuses data only within one Spark application and dies with its executors; writing the result to Parquet makes it readable by every later job, survives failures, and keeps executor memory for execution. Cache inside an application, materialize across them.

solid answer

~50 s

The first question is whether the consumers live inside **one Spark application** or in separately submitted ones. Cached blocks live in that application's executors, so a second `spark-submit` cannot see them at all — `persist()` is simply not an option there. Within one application, persist is right when the dataset is read a handful of times, fits comfortably beside execution memory, and the recomputation it avoids is genuinely expensive. Writing to Parquet on object storage costs a serialization pass and IO, but buys reuse across applications, durability across executor loss, column pruning and predicate pushdown on every re-read, and zero pressure on the memory pool. It also creates a dataset someone must own: a path, a retention policy, a schema. My rule of thumb: cache for tactical reuse inside a job, materialize for anything another team or another schedule will read. `checkpoint()` sits between them — it truncates lineage to reliable storage without becoming a published table.

code

python · 6 lines
python
# tactical: several actions in ONE application
enriched = orders.join(customers, "customer_id")
enriched.persist()
enriched.groupBy("region").sum("amount").write.save(regions_path)
enriched.groupBy("order_date").count().write.save(daily_path)
enriched.unpersist()

go deeper

for a junior

Know that caching helps only when data is read more than once, and that writing results to storage is how another job gets to read them.

for a middle

Explain that cached blocks live in the application's own executors, and compare the costs: memory taken from execution versus a serialization and write pass plus a re-readable columnar file.

for a senior

Bring the operational angle: durability across executor loss, column pruning on re-read, the optimization barrier a cache introduces, and when checkpoint is the right tool for plan depth rather than reuse.

for a principal

Own the tradeoff at pipeline scale — decide which intermediates become owned, scheduled datasets with contracts and which stay tactical caches, and be able to argue the storage cost and lineage-governance consequences either way.

## Reframe the question: what is the scope of reuse? "Cache or write?" is usually asked as a performance question. It is really a **scope and ownership** question, and the performance answer follows from it. - **Within one Spark application**, cached blocks in executor storage memory are reachable by every subsequent action in that application. Reuse is nearly free. - **Across applications**, they are not reachable at all. Each `spark-submit` gets its own driver and executors; when the application ends, its executors and everything they cached are gone. If your "five downstream jobs" are five scheduled submissions, `persist()` is not a slow option — it is a non-option. So establish first whether the five consumers are five actions in one driver, or five jobs on a scheduler. ## What persist gives you Within one application, persisting an expensive intermediate avoids re-running its lineage for every action. The conditions under which it pays: - The dataset is genuinely read **more than once**. A single linear read-transform-write pipeline gains nothing. - The upstream work is expensive — a shuffle, a join, a heavy UDF, a wide scan. Caching a plain file scan often saves less than it costs. - It **fits** beside execution memory. Cached blocks occupy the storage side of the same unified pool that shuffles and joins draw from; a large cache buys reuse and pays for it in spill. - You control the lifetime. In long-running drivers — notebooks, Thrift servers, streaming applications — forgetting `unpersist()` starves execution memory over hours. And one non-obvious cost: caching creates an optimization barrier. Filters and column pruning cannot be pushed through a cached plan into the source scan, so a `cache()` placed too early in a plan can make the job read substantially more data than the uncached version would. ## What writing gives you Materializing to Parquet (or an open table format) on durable storage costs one serialization and write pass. In exchange: - **Cross-application reuse.** Any job, any schedule, any team can read it. - **Durability.** Executor loss, driver restart and job failure do not destroy it; a retry resumes from the materialized dataset instead of recomputing the pipeline. - **Read efficiency on re-read.** Columnar layout means downstream readers prune columns and push filters, so a consumer that needs three of forty columns reads three. A cached DataFrame is reused whole. - **No memory pressure.** The unified pool stays available for execution. - **Debuggability.** You can query the intermediate, diff it across runs, and reason about it after the fact. The costs are equally real: the write pass itself, storage cost, and — the one teams underestimate — **lifecycle ownership**. A materialized intermediate is a dataset with a path, a schema, a retention policy and a consumer contract. Ten of these accumulated without governance become a warehouse of orphaned staging tables nobody dares delete. ## Where checkpoint() sits `DataFrame.checkpoint()` writes to reliable storage *and* truncates lineage, so the plan no longer references the upstream chain. That matters for iterative algorithms whose plans grow with each iteration until the driver spends its time planning rather than executing, and for very deep lineages where recovery from executor loss would be prohibitive. It is a *mechanism*, not a published dataset: the files are managed by Spark under `spark.sparkContext.setCheckpointDir` and are not something downstream jobs should read. `localCheckpoint()` trades reliability for speed by keeping the truncation on executor local disk — useful for plan-depth control, useless for fault tolerance. ## A decision framework I would actually apply 1. **Consumers in another application or another schedule?** → Write it. Nothing else works. 2. **Consumers in this application, read two or three times, dataset small relative to executor memory, upstream expensive?** → Persist, and unpersist when done. 3. **Consumers in this application but the dataset is large, or memory is already tight, or the stage is already spilling?** → Write it and re-read. Being able to prune columns on the re-read often makes this *faster* than caching, not merely safer. 4. **Plan depth or iteration count is the problem, not reuse?** → Checkpoint. 5. **Read exactly once?** → Do neither. This is the most common unnecessary `cache()` in real codebases. ## The organizational layer At scale, the interesting failure is not a badly chosen storage level; it is a pipeline where the same expensive join is recomputed by six teams because none of them knew the others needed it. That is solved by publishing the intermediate as an owned dataset with a schedule and a contract, not by tuning anyone's cache. Conversely, the opposite failure — every intermediate materialized "just in case" — produces storage cost and a lineage nobody can trace. The judgment being tested is whether you can tell which of those two failure modes your organization is currently in.

  • Can two separately submitted Spark applications share a persisted DataFrame?
    No. Cached blocks live in the storage memory of the application's own executors, which are torn down when it ends, and there is no cross-application block registry. The external shuffle service preserves shuffle files, not cache blocks. Sharing across applications requires materializing to durable storage, or an external caching layer that is not part of Spark's memory manager.
  • How does checkpoint() differ from persist() for an iterative algorithm?
    persist keeps the data but leaves the logical plan intact, so each iteration's plan still references everything before it and planning cost grows without bound. checkpoint writes to reliable storage and truncates the lineage, so the next iteration starts from a short plan. For iterative work the plan-depth problem is often the real bottleneck, and only checkpointing addresses it.
  • When is re-reading a written Parquet file actually faster than reading a cached DataFrame?
    When downstream consumers need only a subset of columns or rows. A Parquet re-read prunes columns and pushes filters into the scan, so a consumer touching three of forty columns reads three. A cached DataFrame is reused as a whole. Add the memory the cache was denying execution, and the written version can win on wall-clock time as well as on safety.
  • What is the risk of materializing every expensive intermediate to storage?
    You create datasets without owners. Each one needs a path, a schema, a retention policy and a consumer contract, and without governance they accumulate as staging tables nobody can safely delete or explain. Storage cost is the visible price; an untraceable lineage and stale intermediates that quietly feed production reports are the expensive one.

saying these in an interview costs you the question

  • Assumes a cached DataFrame is visible to other Spark applications
  • Treats cache() as free rather than as memory taken from execution
  • Caches datasets that are read exactly once
  • Ignores that caching blocks filter and column pushdown
  • Confuses checkpoint with a published, readable intermediate dataset

context