ML Engineer MasterClass (October) | 4 seats left

Spark · Partition pruning and column selection
AmazonAmazon Analytics
Amazon · Read and Write Data

Partition pruning and column selection

Inspect directory partition filters and the columns requested by a Parquet scan.

Step 1 of 6 · Learn

Filter a directory partition

partitionBy("category") writes category directories; this is file layout, not Window.partitionBy or an in-memory repartition call. Filtering the partition field lets a reader skip irrelevant directories. Inspect PartitionFilters in the printed physical plan. The temporary dataset is tiny, so it illustrates the mechanism rather than demonstrating a meaningful speedup.

Lesson reference: PySpark and Scala

Filter a directory partition

partitionBy("category") writes category directories; this is file layout, not Window.partitionBy or an in-memory repartition call. Filtering the partition field lets a reader skip irrelevant directories. Inspect PartitionFilters in the printed physical plan. The temporary dataset is tiny, so it illustrates the mechanism rather than demonstrating a meaningful speedup.

PySpark example

from tempfile import TemporaryDirectory
from pyspark.sql.functions import col, count, lit

with TemporaryDirectory(prefix="spark-storage-") as directory:
    path = directory + "/data"
    compact_path = directory + "/compact"
    products.write.partitionBy("category").parquet(path)
    loaded = spark.read.parquet(path).filter(col("category") == "Electronics").select("product_id", "product_name", "category").orderBy("product_id")
    loaded.explain("formatted")
    # Detach only this tiny teaching result before removing its files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema).orderBy("product_id")
    result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val directory = Files.createTempDirectory("spark-storage-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val compact_path = directory.resolve("compact").toString
val result = try {
    products.write.partitionBy("category").parquet(path)
    val loaded = spark.read.parquet(path).filter(col("category") === "Electronics").select("product_id", "product_name", "category").orderBy("product_id")
    loaded.explain("formatted")
    val output = spark.createDataFrame(spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema).orderBy("product_id")
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only this example’s unique temporary directory.
}

Read only required columns

Parquet is columnar. Selecting fewer columns can reduce the columns the scan reads; inspect ReadSchema rather than assuming SELECT-like projection always reads everything. Columns used by filters may still be needed even if omitted from output. This example has no filter. Metadata listing and other overhead remain.

PySpark example

from tempfile import TemporaryDirectory
from pyspark.sql.functions import col, count, lit

with TemporaryDirectory(prefix="spark-storage-") as directory:
    path = directory + "/data"
    compact_path = directory + "/compact"
    products.write.parquet(path)
    loaded = spark.read.parquet(path).select("product_id", "product_name").orderBy("product_id")
    loaded.explain("formatted")
    # Detach only this tiny teaching result before removing its files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema).orderBy("product_id")
    result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val directory = Files.createTempDirectory("spark-storage-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val compact_path = directory.resolve("compact").toString
val result = try {
    products.write.parquet(path)
    val loaded = spark.read.parquet(path).select("product_id", "product_name").orderBy("product_id")
    loaded.explain("formatted")
    val output = spark.createDataFrame(spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema).orderBy("product_id")
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only this example’s unique temporary directory.
}

Combine pruning and projection

A partition filter selects relevant directories, while projection selects required data columns. These are separate optimizations. category can be used to choose the directory without appearing in the result. Inspect both PartitionFilters and ReadSchema in the live plan; exact plan formatting can vary across Spark versions.

PySpark example

from tempfile import TemporaryDirectory
from pyspark.sql.functions import col, count, lit

with TemporaryDirectory(prefix="spark-storage-") as directory:
    path = directory + "/data"
    compact_path = directory + "/compact"
    products.write.partitionBy("category").parquet(path)
    loaded = spark.read.parquet(path).filter(col("category") == "Home").select("product_id", "product_name", "category").orderBy("product_id")
    loaded.explain("formatted")
    # Detach only this tiny teaching result before removing its files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema).orderBy("product_id")
    result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val directory = Files.createTempDirectory("spark-storage-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val compact_path = directory.resolve("compact").toString
val result = try {
    products.write.partitionBy("category").parquet(path)
    val loaded = spark.read.parquet(path).filter(col("category") === "Home").select("product_id", "product_name", "category").orderBy("product_id")
    loaded.explain("formatted")
    val output = spark.createDataFrame(spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema).orderBy("product_id")
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only this example’s unique temporary directory.
}
example.pyPySpark
1. Directory layoutWrite a private Parquet dataset; category partitions are used where shown.
2. Filter and projectInspect PartitionFilters and ReadSchema in the actual plan.
3. Return requested recordsDetach only these few rows before cleanup.
Source products6 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
202Laptop standOffice
203Desk lampHome
204USB-C hubElectronics
205Travel backpackTravel
206Notebook setOffice
Filter a directory partition
Result2 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
204USB-C hubElectronics
The example selects Electronics products 201 and 204. Directory discovery supplies category as a partition column.

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.