ML Engineer MasterClass (October) | 4 seats left

Spark · Cache and persist
AmazonAmazon Analytics
Amazon · Tune Spark Workloads

Cache and persist

Reuse a DataFrame across actions and release storage deliberately.

Step 1 of 6 · Learn

Materialize before reuse

cache marks a DataFrame for reuse. The first action computes and stores partitions; this example uses count, then show. Caching has overhead and is usually unnecessary for six rows. It is valuable when expensive intermediate data is reused. finally releases this example’s cache; the studio’s later validation collection may recompute the result.

Lesson reference: PySpark and Scala

Materialize before reuse

cache marks a DataFrame for reuse. The first action computes and stores partitions; this example uses count, then show. Caching has overhead and is usually unnecessary for six rows. It is valuable when expensive intermediate data is reused. finally releases this example’s cache; the studio’s later validation collection may recompute the result.

PySpark example

from pyspark import StorageLevel
from pyspark.sql.functions import col, sum

cached = orders.filter(col("total") >= 100).cache()
try:
    cached.count()  # Materialize the cache with an action.
    result = cached.select("order_id", "total").orderBy("order_id")
    result.show(truncate=False)
finally:
    cached.unpersist(blocking=True)

Scala example

import org.apache.spark.storage.StorageLevel
import org.apache.spark.sql.functions._

val cached = orders.filter(col("total") >= 100).cache()
val result = try {
    cached.count() // Materialize the cache with an action.
    val output = cached.select("order_id", "total").orderBy("order_id")
    output.show(false)
    output
} finally {
    cached.unpersist(true)
}

Choose a storage level

persist accepts a storage level; DISK_ONLY stores partitions on disk rather than keeping them in memory. It trades memory pressure for disk I/O and is not inherently faster. Choose a level from workload measurements. This example creates a fresh DataFrame and releases it afterward, rather than changing a level on an already-persisted object.

PySpark example

from pyspark import StorageLevel
from pyspark.sql.functions import col, sum

cached = orders.filter(col("total") >= 100).cache()
try:
    cached.count()  # Materialize the cache with an action.
    result = cached.select("order_id", "total").orderBy("order_id")
    result.show(truncate=False)
finally:
    cached.unpersist(blocking=True)

Scala example

import org.apache.spark.storage.StorageLevel
import org.apache.spark.sql.functions._

val cached = orders.filter(col("total") >= 100).cache()
val result = try {
    cached.count() // Materialize the cache with an action.
    val output = cached.select("order_id", "total").orderBy("order_id")
    output.show(false)
    output
} finally {
    cached.unpersist(true)
}

Reuse a filtered input for aggregation

A reusable intermediate can feed several actions. Here count materializes the filtered input, then a grouped sum reads it. Only the filtered input is persisted; downstream aggregation and sorting still do work. unpersist removes stored blocks, not the DataFrame definition. Later actions can recompute its lineage.

PySpark example

from pyspark import StorageLevel
from pyspark.sql.functions import col, sum

cached = orders.filter(col("total") >= 100).cache()
try:
    cached.count()  # Materialize the cache with an action.
    result = cached.groupBy("status").agg(sum("total").alias("total_value")).orderBy("status")
    result.show(truncate=False)
finally:
    cached.unpersist(blocking=True)

Scala example

import org.apache.spark.storage.StorageLevel
import org.apache.spark.sql.functions._

val cached = orders.filter(col("total") >= 100).cache()
val result = try {
    cached.count() // Materialize the cache with an action.
    val output = cached.groupBy("status").agg(sum("total").alias("total_value")).orderBy("status")
    output.show(false)
    output
} finally {
    cached.unpersist(true)
}
example.pyPySpark
1. Mark for reusecache or persist selects storage behavior; it does not eagerly load all rows.
2. Materialize and reusecount triggers computation; a later action can reuse stored partitions.
3. Release in finallyunpersist frees this cache even if the example fails. The studio may recompute the returned DataFrame for validation.
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Materialize before reuse
Result3 rows
order_idtotal
1002149
1004219.99
1006120
The example reuses the three orders worth at least 100, then releases their stored partitions.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.