ML Engineer MasterClass (October) | 4 seats left

Spark · Small files
AmazonAmazon Analytics
Amazon · Read and Write Data

Small files

Control tiny teaching writes and understand why production file sizing needs measurement.

Step 1 of 6 · Learn

Reduce writer partitions deliberately

Many small files add listing, opening, and scheduling overhead. Coalescing a tiny dataset to one writer can reduce files, but coalesce(1) is not a general production recommendation: it limits write parallelism and can bottleneck large data. File counts also depend on directory partitioning and writer options. This example checks preserved records, not an assumed universal file count.

Lesson reference: PySpark and Scala

Reduce writer partitions deliberately

Many small files add listing, opening, and scheduling overhead. Coalescing a tiny dataset to one writer can reduce files, but coalesce(1) is not a general production recommendation: it limits write parallelism and can bottleneck large data. File counts also depend on directory partitioning and writer options. This example checks preserved records, not an assumed universal file count.

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.repartition(2).write.parquet(path)
    loaded = spark.read.parquet(path).select("product_id", "product_name").orderBy("product_id")
    # 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.repartition(2).write.parquet(path)
    val loaded = spark.read.parquet(path).select("product_id", "product_name").orderBy("product_id")
    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.
}

Compact into a new destination

Repeated appends can create small files. This example writes non-overlapping batches, reads all of them, and compacts into a different path. Reading from a path while overwriting it is unsafe without a suitable storage/table protocol. Each nonempty batch uses one writer and disables record caps, making file counts controlled for this tiny exercise.

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.filter((col("product_id") >= 201) & (col("product_id") < 203)).coalesce(1).write.option("maxRecordsPerFile", 0).mode("append").parquet(path)
    products.filter((col("product_id") >= 203) & (col("product_id") < 205)).coalesce(1).write.option("maxRecordsPerFile", 0).mode("append").parquet(path)
    products.filter((col("product_id") >= 205) & (col("product_id") < 207)).coalesce(1).write.option("maxRecordsPerFile", 0).mode("append").parquet(path)
    source = spark.read.parquet(path)
    source_file_count = len(source.inputFiles())
    source.coalesce(1).write.option("maxRecordsPerFile", 0).parquet(compact_path)
    compacted = spark.read.parquet(compact_path)
    loaded = compacted.agg(count("*").alias("product_count")).select(lit(source_file_count).alias("source_files"), lit(len(compacted.inputFiles())).alias("compacted_files"), col("product_count"))
    # Detach only this tiny teaching result before removing its files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema)
    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.filter((col("product_id") >= 201) && (col("product_id") < 203)).coalesce(1).write.option("maxRecordsPerFile", 0).mode("append").parquet(path)
    products.filter((col("product_id") >= 203) && (col("product_id") < 205)).coalesce(1).write.option("maxRecordsPerFile", 0).mode("append").parquet(path)
    products.filter((col("product_id") >= 205) && (col("product_id") < 207)).coalesce(1).write.option("maxRecordsPerFile", 0).mode("append").parquet(path)
    val source = spark.read.parquet(path)
    val source_file_count = source.inputFiles.length
    source.coalesce(1).write.option("maxRecordsPerFile", 0).parquet(compact_path)
    val compacted = spark.read.parquet(compact_path)
    val loaded = compacted.agg(count("*").alias("product_count")).select(lit(source_file_count).alias("source_files"), lit(compacted.inputFiles.length).alias("compacted_files"), col("product_count"))
    val output = spark.createDataFrame(spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema)
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only this example’s unique temporary directory.
}

Limit records per file

maxRecordsPerFile sets a row ceiling per output file, not a target byte size. Even one writer can create multiple files when the cap is reached. Production choices should consider compressed bytes, row widths, partition layout, and downstream reads. In this controlled example, six rows and one writer make the effect easy to inspect.

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.coalesce(1).write.option("maxRecordsPerFile", 2).parquet(path)
    saved = spark.read.parquet(path)
    loaded = saved.agg(count("*").alias("product_count")).select(lit(len(saved.inputFiles())).alias("data_files"), col("product_count"))
    # Detach only this tiny teaching result before removing its files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema)
    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.coalesce(1).write.option("maxRecordsPerFile", 2).parquet(path)
    val saved = spark.read.parquet(path)
    val loaded = saved.agg(count("*").alias("product_count")).select(lit(saved.inputFiles.length).alias("data_files"), col("product_count"))
    val output = spark.createDataFrame(spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema)
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only this example’s unique temporary directory.
}
example.pyPySpark
1. Small input batchesUse explicit writer controls for these six teaching rows.
2. Write or compactUse a separate destination when compacting an existing dataset.
3. Inspect and preserveCheck records and, where requested, actual data-file counts.
Source products6 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
202Laptop standOffice
203Desk lampHome
204USB-C hubElectronics
205Travel backpackTravel
206Notebook setOffice
Reduce writer partitions deliberately
Result6 rows
product_idproduct_name
201Wireless keyboard
202Laptop stand
203Desk lamp
204USB-C hub
205Travel backpack
206Notebook set
The six product records must survive either physical writing strategy.

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.