Partitioned Table Writes

partitionBy writes one key=value directory per value, and a filter on that column reads only matching directories. Partition by a low-cardinality column that queries filter on, usually a date.

Writing order lines partitioned by month, then pruning and overwritingJavaScript
import contextlib, io, os, re
path = "out/lines_by_month"
monthly = lines.withColumn("month", F.date_format("order_ts", "yyyy-MM"))
monthly.write.mode("overwrite").partitionBy("month").parquet(path)
def show_layout():
    dirs = sorted(d for d in os.listdir(path) if d.startswith("month="))
    files = sum(len([f for f in os.listdir(f"{path}/{d}") if f.endswith(".parquet")])
                for d in dirs)
    print(f"{len(dirs)} directories, {files} files: {dirs[0]} ... {dirs[-1]}")
show_layout()
june = spark.read.parquet(path).where("month = '2026-06'")
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
    june.explain()
print(f"{june.count():,} rows;", re.search(r"PartitionFilters: \[[^\]]*\]", buf.getvalue())[0])
for mode in ("dynamic", "static"):                # rewrite June only; static is the default
    spark.conf.set("spark.sql.sources.partitionOverwriteMode", mode)
    june_rows = monthly.where("month = '2026-06'")
    june_rows.write.mode("overwrite").partitionBy("month").parquet(path)
    print(f"{mode}: ", end="")
    show_layout()
Output
18 directories, 21 files: month=2025-01 ... month=2026-06
75,867 rows; PartitionFilters: [isnotnull(month#44), (month#44 = 2026-06)]
dynamic: 18 directories, 21 files: month=2025-01 ... month=2026-06
static: 1 directories, 1 files: month=2026-06 ... month=2026-06

Only 21 files, not 4 x 18: the input is sorted by time, so each task touched few months. PartitionFilters prune whole directories before any file opens. The overwrite shows the classic trap: in dynamic mode June's rows replaced only month=2026-06, while in static mode, the default, the same call deleted the other 17 months. Use dynamic for jobs that rebuild single partitions.