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.
}