PySpark courseLesson 10 of 10
PySpark course · Lesson 10 of 10
Writing Efficient Output: File Sizes, Small Files and Partitioned Writes
Control how many files PySpark writes and how big they are: repartition vs coalesce before a write, partitioned writes, the small files problem, fewer shuffles and custom partitioners.
On this page
- Sample data and helpers
- Repartition before a write
- What it is
- Partitioned writes: the multiplication problem
- Too much data per folder: spread it, or cap file size
- Range partitioning and sorting for better skipping
- Pitfalls
- In interviews
- Coalesce on output
- What it is
- When coalesce is the right call
- Pitfalls and interview angle
- Optimising file sizes
- What “good” looks like
- How to aim for a target size
- Pitfalls
- In interviews
- The small files problem
- What it is and why it hurts
- How Spark packs small files on read
- Fixes
- In interviews
- Avoiding unnecessary wide transformations
- Narrow vs wide
- Combine work that shares a key
- Prefer operations that combine before shuffling
- A checklist
- In interviews
- Custom partitioners
- What a partitioner is
- Custom partitioners on RDDs
- DataFrames: no custom partitioner, but hashed expressions
- Pitfalls and interview angle
- Practice questions
- Key takeaways
A Spark job that computes the right answer can still leave a mess behind: ten thousand 50 KB files, one 40 GB file, or a folder per customer. The shape of the output decides how fast every later reader runs, and the shuffles you introduce to shape it decide how fast the job itself runs. This lesson shows how a write turns partitions into files, how to control file count and size, why small files are a problem, how to avoid unnecessary wide transformations, and what custom partitioners do.
Sample data and helpers
400,000 synthetic events across 4 countries and 3 dates, created in 8 partitions. summary() counts the Parquet files under a folder and reports the smallest and largest sizes.
import os, glob, tempfile, re, io, contextlib
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder.master("local[4]").appName("output")
.config("spark.sql.shuffle.partitions", "16")
.config("spark.sql.adaptive.enabled", "false") # predictable partition counts for the demo
.getOrCreate())
spark.sparkContext.setLogLevel("ERROR")
base = tempfile.mkdtemp()
events = spark.range(0, 400_000, numPartitions=8).select(
F.col("id").alias("event_id"),
F.element_at(F.array(*[F.lit(c) for c in ["IN", "UK", "US", "DE"]]),
(F.col("id") % 4 + 1).cast("int")).alias("country"),
F.date_add(F.lit("2026-01-01").cast("date"), (F.col("id") % 3).cast("int")).alias("event_date"),
(F.col("id") % 997 * 0.1).alias("amount"),
)
def summary(path):
sizes = sorted(os.path.getsize(f) for f in glob.glob(path + "/**/*.parquet", recursive=True))
return f"{len(sizes)} files, smallest {sizes[0] / 1024:.1f} KiB, largest {sizes[-1] / 1024:.1f} KiB"
def plan(df):
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
df.explain()
text = re.sub(r"#\d+L?", "", buf.getvalue())
print(re.sub(r", \[plan_id=\d+\]", "", text).strip())
print(events.rdd.getNumPartitions())
8
The rule behind everything below: each task writes its own files. A write of a DataFrame with N partitions produces up to N files (empty partitions write nothing), and with partitionBy, up to N files per output folder.
Repartition before a write
What it is
repartition(n) or repartition(n, cols) adds a shuffle that redistributes rows into exactly n partitions, either round-robin (no columns) or by hashing the columns. Placed just before a write, it sets the number of output files and makes them roughly equal in size. With columns and no n, it uses spark.sql.shuffle.partitions.
Partitioned writes: the multiplication problem
partitionBy("event_date") on the writer creates one folder per date. Without preparing the data, every task writes a file into every folder it has rows for:
p = base + "/by_date_naive"
events.write.mode("overwrite").partitionBy("event_date").parquet(p)
print(summary(p))
p2 = base + "/by_date_repart"
events.repartition("event_date").write.mode("overwrite").partitionBy("event_date").parquet(p2)
print(summary(p2))
print(sorted(d for d in os.listdir(p2) if d.startswith("event_date")))
24 files, smallest 78.4 KiB, largest 78.5 KiB
3 files, smallest 585.9 KiB, largest 585.9 KiB
['event_date=2026-01-01', 'event_date=2026-01-02', 'event_date=2026-01-03']
8 input partitions × 3 dates = 24 files. On a real job with 2,000 tasks and 365 dates that becomes hundreds of thousands of small files. repartition("event_date") sends all rows of one date to one task, giving exactly one file per folder.
Too much data per folder: spread it, or cap file size
One file per folder is ideal only if a folder’s data is a sensible file size. When a date holds tens of gigabytes, one task would write one enormous file slowly. Two options:
p3 = base + "/by_date_spread"
(events.repartition(6, "event_date", F.col("event_id") % 2)
.write.mode("overwrite").partitionBy("event_date").parquet(p3))
print(summary(p3))
p4 = base + "/capped"
(events.repartition("event_date").write.mode("overwrite")
.option("maxRecordsPerFile", 50_000).partitionBy("event_date").parquet(p4))
print(summary(p4))
5 files, smallest 297.7 KiB, largest 585.9 KiB
9 files, smallest 152.3 KiB, largest 225.1 KiB
- Repartition by the folder column plus a spreading expression (here
event_id % 2) lets a few tasks share each folder. Hash collisions mean the split is not perfectly even: one date got two files and others one. A random salt (F.floor(F.rand() * k)) is common, but it is non-deterministic across retries, so a hash of a key is safer. maxRecordsPerFile(option orspark.sql.files.maxRecordsPerFile, default 0 = unlimited) makes a task start a new file after N rows. It bounds file size without another shuffle but cannot merge small files from different tasks.
Range partitioning and sorting for better skipping
repartitionByRange(n, col) samples the data to choose boundaries, so each partition holds a contiguous range of values. Combined with sortWithinPartitions, every output file covers a narrow, non-overlapping range, which makes Parquet min/max statistics (and table-format data skipping) effective:
r = events.repartitionByRange(4, "event_id").sortWithinPartitions("event_id")
r.write.mode("overwrite").parquet(base + "/ranged")
stats = (spark.read.parquet(base + "/ranged")
.groupBy(F.input_file_name().alias("f"))
.agg(F.min("event_id").alias("min_id"), F.max("event_id").alias("max_id"),
F.count("*").alias("rows"))
.orderBy("min_id").drop("f"))
stats.show()
+------+------+------+
|min_id|max_id| rows|
+------+------+------+
| 0| 99822| 99823|
| 99823|199899|100077|
|199900|299238| 99339|
|299239|399999|100761|
+------+------+------+
The ranges do not overlap, so a query for event_id = 250000 can skip three of the four files from their footers alone. Boundaries come from a sample, so partitions are close to, not exactly, equal.
Pitfalls
- Repartitioning costs a full shuffle. Do it once, at the end, not between every step.
- Repartitioning by a skewed column sends a huge share of rows to one task; add a spreading expression or cap records per file.
partitionByon a high-cardinality column (user ID, timestamp) creates a folder per value. No repartitioning can fix that; pick a coarser partition column.- AQE (on by default) may coalesce shuffle partitions before the write, changing the file count you expected from
repartition(n)by column; aREPARTITIONwith an explicit number is respected, but check the output.
In interviews
“Your job writes 200,000 small files to a date-partitioned table. Why, and how do you fix it?” Strong answer: every task writes to every partition folder it touches; repartition by the partition column before writing (plus a spreading key or maxRecordsPerFile for big partitions), and compact existing files.
Coalesce on output
What it is
coalesce(n) reduces the number of partitions without a shuffle by merging existing partitions on the same executors. It is cheap, so it is tempting as “write fewer files”. Two facts make it a trap:
- It can only decrease the partition count (
coalesce(100)on 10 partitions does nothing). - Because it is not a shuffle boundary, the reduced parallelism applies to the whole stage before it, back to the previous shuffle or the read.
heavy = events.withColumn("score", F.sha2(F.col("event_id").cast("string"), 256))
plan(heavy.coalesce(2))
plan(heavy.repartition(2))
== Physical Plan ==
Coalesce 2
+- *(1) Project [id AS event_id, element_at([IN,UK,US,DE], cast(((id % 4) + 1) as int), None, true) AS country, date_add(2026-01-01, cast((id % 3) as int)) AS event_date, (cast((id % 997) as double) * 0.1) AS amount, sha2(cast(cast(id as string) as binary), 256) AS score]
+- *(1) Range (0, 400000, step=1, splits=8)
== Physical Plan ==
Exchange RoundRobinPartitioning(2), REPARTITION_BY_NUM
+- *(1) Project [id AS event_id, element_at([IN,UK,US,DE], cast(((id % 4) + 1) as int), None, true) AS country, date_add(2026-01-01, cast((id % 3) as int)) AS event_date, (cast((id % 997) as double) * 0.1) AS amount, sha2(cast(cast(id as string) as binary), 256) AS score]
+- *(1) Range (0, 400000, step=1, splits=8)
With Coalesce 2, the hashing (sha2) and everything in stage 1 run in 2 tasks instead of 8. With repartition(2) there is an Exchange: stage 1 still runs with 8 tasks, then a shuffle feeds 2 writing tasks. On a cluster, coalesce(1) after an expensive transformation can turn a 10-minute job into a 3-hour job on a single core.
When coalesce is the right call
agg = events.groupBy("country", "event_date").agg(F.sum("amount").alias("revenue"))
print(agg.rdd.getNumPartitions())
agg.write.mode("overwrite").parquet(base + "/agg_default")
print(summary(base + "/agg_default"))
agg.coalesce(1).write.mode("overwrite").parquet(base + "/agg_one")
print(summary(base + "/agg_one"))
16
9 files, smallest 1.0 KiB, largest 1.0 KiB
1 files, smallest 1.1 KiB, largest 1.1 KiB
After an aggregation the data is small (12 rows here) but spread over 16 shuffle partitions, so the default write produced 9 tiny files (empty partitions wrote none). Coalescing a small, already computed result is fine: the expensive work happened in the stage before the shuffle, at full parallelism. (With AQE on, Spark would usually coalesce these tiny shuffle partitions automatically.)
coalesce(n) |
repartition(n) |
|
|---|---|---|
| Shuffle | No | Yes |
| Can increase partitions | No | Yes |
| Resulting sizes | Uneven (merges neighbours) | Even (round-robin or hash) |
| Effect on upstream | Reduces parallelism of the whole preceding stage | Upstream keeps its parallelism |
| Use before a write when | The data is already small or reduced, and the preceding stage is cheap | You need fewer, even files after heavy work, or more partitions |
Pitfalls and interview angle
coalesce(1)to “get one CSV file” is a common anti-pattern on big data; write normally and merge outside Spark if a single file is truly required.- Interviewers ask “repartition vs coalesce?” constantly. Mention the shuffle, decrease-only, uneven sizes and, the point most candidates miss, that coalesce lowers parallelism of the upstream stage.
Optimising file sizes
What “good” looks like
There is no universal number, but the common guidance is files in the hundreds of megabytes up to around 1 GB for analytical tables on object storage: big enough that per-file overhead (open, footer read, listing, scheduling a task) is small relative to the data, small enough to keep many tasks busy and to rewrite cheaply. Delta Lake’s OPTIMIZE, for example, targets about 1 GB per file by default. Check your own engine’s or platform’s guidance; managed platforms often auto-tune file sizes.
How to aim for a target size
Spark does not have a “target file size” option for plain writes. You steer it by the number of rows per task:
- Estimate the output size (from a previous run or
spark.read.parquet(path)input size and compression ratio). - Choose the partition count as
output size / target file size, andrepartition(n)(or by the partition column with a spread factor) before the write. - Cap with
maxRecordsPerFileas a guard against skewed partitions. - Let AQE help: with AQE on,
spark.sql.adaptive.advisoryPartitionSizeInBytes(64 MB by default) guides how shuffle partitions are coalesced, which also shapes the files written after a shuffle. - Sort within partitions by commonly filtered columns: it improves compression and min/max skipping, often shrinking files noticeably.
Codec and format matter too. Parquet uses Snappy by default; zstd usually compresses better at some CPU cost (spark.sql.parquet.compression.codec).
Pitfalls
- Optimising file size for the writer and forgetting the reader: many readers filter by date, so files must be well-sized per partition folder, not overall.
- Very large files (multiple GB) limit read parallelism for formats or codecs that cannot be split and make rewrites (updates, deletes, compaction) expensive.
- Measuring sizes in the Spark UI’s shuffle metrics rather than on disk: compression changes the picture.
In interviews
“How would you control output file size?” Answer with the rows-per-task model, repartition counts derived from estimated size, maxRecordsPerFile, AQE advisory size, sorting, and table-format compaction (OPTIMIZE) or auto-compaction where available.
The small files problem
What it is and why it hurts
A table made of many tiny files reads slowly, even when the total data is small:
- Listing and metadata: object stores list files page by page; thousands of files mean thousands of requests before a single byte is read.
- Per-file overhead: each file needs an open, a footer read and statistics parsing. For a 30 KB file that overhead dominates.
- Scheduling: one task per file (or per small group of files) adds scheduler overhead and too many tiny tasks.
- Worse compression and statistics: small row groups compress poorly and skip poorly.
Typical causes: streaming micro-batches writing every few seconds, high-cardinality partitionBy, many writing tasks into each partition folder, frequent small appends and over-partitioned shuffles.
How Spark packs small files on read
small = base + "/small_files"
events.repartition(200).write.mode("overwrite").parquet(small)
print(summary(small))
print(spark.read.parquet(small).rdd.getNumPartitions())
spark.conf.set("spark.sql.files.openCostInBytes", str(64 * 1024 * 1024))
print(spark.read.parquet(small).rdd.getNumPartitions())
spark.conf.set("spark.sql.files.openCostInBytes", str(4 * 1024 * 1024))
200 files, smallest 17.2 KiB, largest 17.4 KiB
7
100
Spark combines small files into read partitions up to spark.sql.files.maxPartitionBytes (128 MB by default), but charges each file an assumed “open cost” (spark.sql.files.openCostInBytes, 4 MB by default) so that it does not pack too many files into one task. Raising the open cost to 64 MB made Spark treat each file as expensive and produced 100 read tasks for the same 200 files. Tuning these settings helps the scheduler, but the listing and per-file I/O cost remain: the real fix is fewer files.
Fixes
- Prevent at write time: repartition by the partition column, size partitions sensibly, avoid high-cardinality partition columns.
- Compact: read the small files and rewrite them as fewer, larger files, then swap them in.
compacted = base + "/compacted"
spark.read.parquet(small).repartition(2).write.mode("overwrite").parquet(compacted)
print(len(glob.glob(compacted + "/*.parquet")), "files,", spark.read.parquet(compacted).count(), "rows")
2 files, 400000 rows
Each compacted file is about 1.3 MB here, versus 17 KB before.
Compacting plain Parquet in place is unsafe: readers may see a half-replaced folder, and a failure leaves duplicates or gaps. Write to a new location and switch a table pointer, or use a table format: Delta Lake’s OPTIMIZE, Iceberg’s rewrite_data_files and similar commands compact transactionally (see Delta Lake transactions).
- For streaming: use longer trigger intervals, write to a table format and schedule compaction, or use features such as optimised writes and auto-compaction where your platform offers them.
In interviews
“What is the small files problem and how do you solve it?” Strong answers name the causes (streaming, over-partitioning, many writers per folder), the costs (listing, per-file overhead, scheduling, poor compression), and both prevention and transactional compaction.
Avoiding unnecessary wide transformations
Narrow vs wide
A narrow transformation computes each output partition from one input partition (select, filter, withColumn, union, coalesce). A wide transformation needs rows from many input partitions and therefore a shuffle (groupBy, join without broadcast, distinct, orderBy, repartition, windows with partitionBy). Shuffles write data to disk, send it over the network and start a new stage, so the fewest, smallest shuffles usually win. See partitions, shuffles and skew.
Narrow operations chain into one stage with no Exchange:
plan(events.filter("country = 'IN'").select("event_id", "amount"))
== Physical Plan ==
*(1) Project [id AS event_id, (cast((id % 997) as double) * 0.1) AS amount]
+- *(1) Filter (element_at([IN,UK,US,DE], cast(((id % 4) + 1) as int), None, true) = IN)
+- *(1) Range (0, 400000, step=1, splits=8)
Combine work that shares a key
Two aggregations over the same key, joined back together, read the data twice and shuffle twice (plus a join):
plan(events.groupBy("country").count().join(events.groupBy("country").agg(F.avg("amount")), "country"))
== Physical Plan ==
*(4) Project [country, count, avg(amount)]
+- *(4) BroadcastHashJoin [country], [country], Inner, BuildRight, false, false
:- *(4) HashAggregate(keys=[country], functions=[count(1)])
: +- Exchange hashpartitioning(country, 16), ENSURE_REQUIREMENTS
: +- *(1) HashAggregate(keys=[country], functions=[partial_count(1)])
: +- *(1) Project [element_at([IN,UK,US,DE], cast(((id % 4) + 1) as int), None, true) AS country]
: +- *(1) Range (0, 400000, step=1, splits=8)
+- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, false]),false)
+- *(3) HashAggregate(keys=[country], functions=[avg(amount)])
+- Exchange hashpartitioning(country, 16), ENSURE_REQUIREMENTS
+- *(2) HashAggregate(keys=[country], functions=[partial_avg(amount)])
+- *(2) Project [element_at([IN,UK,US,DE], cast(((id % 4) + 1) as int), None, true) AS country, (cast((id % 997) as double) * 0.1) AS amount]
+- *(2) Range (0, 400000, step=1, splits=8)
One agg with both functions needs a single scan and a single shuffle:
plan(events.groupBy("country").agg(F.count("*").alias("n"), F.avg("amount").alias("avg")))
== Physical Plan ==
*(2) HashAggregate(keys=[country], functions=[count(1), avg(amount)])
+- Exchange hashpartitioning(country, 16), ENSURE_REQUIREMENTS
+- *(1) HashAggregate(keys=[country], functions=[partial_count(1), partial_avg(amount)])
+- *(1) Project [element_at([IN,UK,US,DE], cast(((id % 4) + 1) as int), None, true) AS country, (cast((id % 997) as double) * 0.1) AS amount]
+- *(1) Range (0, 400000, step=1, splits=8)
Likewise, a window over a key can often replace “aggregate, then join back to the detail rows”.
Prefer operations that combine before shuffling
On RDDs, reduceByKey combines values inside each partition before the shuffle; groupByKey ships every value across the network and builds the full list per key in memory. Both give the same answer here, but at scale only one survives skewed keys:
pairs = spark.sparkContext.parallelize([("IN", 1), ("UK", 1), ("IN", 1), ("US", 1), ("IN", 1)], 2)
print(sorted(pairs.reduceByKey(lambda a, b: a + b).collect()))
print(sorted((k, sum(v)) for k, v in pairs.groupByKey().collect()))
[('IN', 3), ('UK', 1), ('US', 1)]
[('IN', 3), ('UK', 1), ('US', 1)]
DataFrame groupBy().agg() already does partial aggregation, as the partial_count step in the plans above shows.
A checklist
- Filter and select early so shuffles carry fewer rows and columns (Spark pushes many filters down automatically, but not through UDFs or some joins).
- Broadcast small tables instead of shuffling both sides of a join.
- Avoid
distinct()andorderBy()that nobody needs; a global sort is a full shuffle plus range partitioning. Sorting before a write is only useful within partitions (sortWithinPartitions). - Reuse partitioning: consecutive operations on the same key (a join then an aggregation on
customer_id) can share one shuffle when the partitioning matches; an unnecessaryrepartitionin between breaks that. - Cache a DataFrame that several actions reuse, so the shuffles before it are not recomputed.
- Pre-aggregate before joins when the result allows.
In interviews
“How would you reduce shuffles in this job?” Read the plan and count the Exchange nodes, then apply the checklist. Being able to explain why reduceByKey beats groupByKey, and that DataFrame aggregations combine partially, is often probed.
Custom partitioners
What a partitioner is
A partitioner decides which partition each key goes to during a shuffle. The default is a hash partitioner: hash(key) mod n. It spreads keys well on average but gives you no control over which keys share a partition, and a hot key always lands in one partition.
Custom partitioners on RDDs
On pair RDDs, partitionBy(n, partitionFunc) accepts any function that maps a key to an integer; the partition is partitionFunc(key) % n. Use it when you need specific keys together, for example one partition per region for a downstream system, or to isolate a known hot key:
regions = {"IN": 0, "UK": 1, "DE": 1, "US": 2}
def by_region(key):
return regions.get(key, 3)
custom = pairs.partitionBy(4, by_region)
print(custom.glom().map(lambda part: sorted({k for k, _ in part})).collect())
[['IN'], ['UK'], ['US'], []]
Compare with the default hash partitioner on the full event data, where two countries happen to share a partition and one partition stays empty:
keyed = events.select("country", "event_id").rdd.map(lambda r: (r.country, r.event_id))
print(keyed.partitionBy(4).glom().map(lambda p: sorted({k for k, _ in p})).collect())
print(keyed.partitionBy(4, by_region).glom().map(len).collect())
[['IN'], ['DE', 'US'], ['UK'], []]
[100000, 200000, 100000, 0]
The custom function put UK and DE together on purpose (200,000 rows). A partitioner gives control, but you are now responsible for balance: a poor function creates skew just as easily as it prevents it. Converting DataFrames to RDDs also loses Catalyst optimisation and adds Python serialisation, so this is a tool for special cases.
DataFrames: no custom partitioner, but hashed expressions
The DataFrame API has no pluggable partitioner. repartition(n, expr) hashes the expression’s value; it does not use the value as a partition number:
df_custom = events.repartition(4, F.when(F.col("country") == "IN", 0)
.when(F.col("country").isin("UK", "DE"), 1)
.otherwise(2))
print(df_custom.groupBy(F.spark_partition_id().alias("p"))
.agg(F.sort_array(F.collect_set("country")).alias("countries")).orderBy("p").collect())
[Row(p=2, countries=['US']), Row(p=3, countries=['DE', 'IN', 'UK'])]
Values 0 and 1 both hashed to partition 3, so IN, UK and DE ended up together. What DataFrames do offer:
| Need | DataFrame tool |
|---|---|
| Rows with the same key together | repartition(n, key) (hash) |
| Contiguous ranges per partition | repartitionByRange(n, key) |
| Split a hot key across partitions | Salting: add floor(rand() * k) or a hash of another column to the key |
| Control files per output folder | repartition by the folder column plus a spread factor |
| Persisted layout for repeated joins | Bucketing (see joins and join strategy) |
Pitfalls and interview angle
- A custom partitioner function must be deterministic; random output breaks shuffle retries.
- Keep partition counts consistent between RDDs you join: two RDDs partitioned with the same partitioner and count can be joined without another shuffle.
- Interview prompt: “How would you make sure all records for a region are processed by the same task?” Answer: hash repartition by region in DataFrames (or a custom partitioner on RDDs), and discuss balance and hot keys.
Practice questions
A daily job writes a table partitioned by date with 4,000 tasks and produces about 4,000 files per date folder. What is happening and how do you fix it?
Every task holds rows for the date being written and writes its own file into that folder. Repartition by the partition column before writing (df.repartition("date")), adding a spreading expression or maxRecordsPerFile if one date is too big for one file, and compact existing folders.
Why can coalesce(1) before a write make the whole job slower, while repartition(1) does not?
coalesce is not a shuffle boundary, so the single partition applies to the whole stage before it: all upstream transformations run in one task. repartition(1) adds a shuffle, so upstream work runs at full parallelism and only the final write is single-task.
What does spark.sql.files.openCostInBytes do?
It is the estimated cost, in bytes, of opening a file. When Spark packs files into read partitions up to maxPartitionBytes, it adds this cost per file, so many small files are spread across more tasks rather than crammed into one. It changes scheduling, not the underlying per-file I/O.
How would you choose the number of partitions for a write that should produce roughly 512 MB files?
Estimate the compressed output size, divide by 512 MB and repartition to that count (by the partition column with a spread factor when writing partitioned data). Add maxRecordsPerFile as a guard, and adjust from the sizes observed on disk after the first runs.
Why is reduceByKey preferred over groupByKey?
reduceByKey combines values within each partition before shuffling, so it moves one partial result per key per partition. groupByKey shuffles every value and builds the full list per key in memory, which is slower and can fail on hot keys.
You call df.repartition(4, F.col(“region_id”)) where region_id is 0 to 3. Are regions mapped to partitions 0 to 3?
No. DataFrame repartitioning hashes the expression value, so several region IDs may land in the same partition and some partitions may be empty. If you need an exact mapping, use an RDD with a custom partitioner, or accept hash placement and only rely on co-location.
Key takeaways
- Each task writes its own files; with
partitionBy, up to one file per task per folder. - Repartition by the partition column before a partitioned write, and spread or cap (
maxRecordsPerFile) folders that are too large. coalesceavoids a shuffle but lowers the parallelism of the whole preceding stage; use it only on small, already reduced data.- Aim for files in the hundreds of megabytes; steer size through rows per task, sorting and compression.
- Small files cost listing, per-file overhead and scheduling; prevent them at write time and compact transactionally.
- Fewer, smaller shuffles win: combine aggregations, filter early, broadcast small tables, prefer combining operations.
- Custom partitioners exist for RDDs; DataFrames hash expressions, so use range partitioning, salting or bucketing for control.
Progress is saved in this browser only. No account needed.

