Quick summary:
In this composite case study, a nightly PySpark job on Azure Synapse was stalling at 80% while one task ran for two hours and the other 199 finished in minutes. Our team at ScriptsHub Technologies traced the stall to a skewed sort-merge join between 1.2 billion event rows and an 8,000-row account dimension, then replaced it with an explicit PySpark broadcast join. Runtime fell from 3 hours 40 minutes to 38 minutes, Spark pool consumption for the workload dropped by about 60%, and the downstream Power BI refresh landed nearly three hours earlier.
Why is my PySpark join stuck on one task?
A single long-running task during a join usually means shuffle skew: a few join keys hold most of the rows, so one partition does most of the work.
If your Spark UI shows 199 of 200 tasks finishing quickly while one task runs for hours, and your job stalls near the end of the join stage, the likeliest cause is skew in a shuffled PySpark join. This case study follows a SaaS enrichment pipeline on Azure Synapse with these symptoms.
The job read roughly 1.2 billion event rows from an Azure Data Lake Storage Gen2 landing zone, joined them to an 8,000-row account dimension, and wrote Delta output to a curated layer. Power BI refreshes ran two to three hours late, and the Spark pool billed large-executor time the workload did not need. The brief was to fix the stall without redesigning the schema or upsizing the pool.
What is a PySpark broadcast join, and why did this job need one?
A PySpark broadcast join copies a small table to every executor so the large table is joined locally and never shuffled across the network.
Our team opened the physical plan in the Spark web UI and confirmed the query was running as a SortMergeJoin, the default strategy Spark picks when neither side qualifies for broadcast. That strategy shuffles and sorts both sides by the join key. With 1.2 billion rows on the event side, three highly active accounts carrying 41% of all events hashed to the same shuffle partition, which became the straggler.
To confirm it, compare the maximum and median shuffle read size in the stage’s task summary metrics; a maximum many times the median confirms skew.
The dimension should have been broadcast automatically, but the planner’s size estimate for it sat just above the 10 MB default for spark.sql.autoBroadcastJoinThreshold, as documented in Apache’s Spark SQL performance tuning guide. Adaptive query execution (AQE) was enabled, but the hot partition still ran as one task, so its skew-join split never applied. AQE splits a partition only when spark.sql.adaptive.skewJoin.enabled is true and the partition exceeds both five times the median and 256 MB, so check that setting first. Microsoft’s guide to adaptive query execution also notes that AQE may only switch to a broadcast join after both sides have been shuffled, so a statically planned broadcast is usually faster.
Broadcast join vs salting vs AQE tuning: which fix works best for skew?
For a small, stable dimension joined to a skewed fact table, an explicit broadcast join is the lowest-effort and fastest fix.
We tested three PySpark optimization techniques. Enabling and tuning AQE skew-join handling lets Spark split the oversized partition, but every row on the large side still gets shuffled. Salting the hot keys spreads them across partitions reliably, though it requires a matching salt column on the dimension and a more invasive code change. Broadcasting the 8,000-row dimension removes the shuffle of the large side altogether, with no schema change.

Table 1. Runtime, code change, and maintenance cost of three skew fixes on the case study workload. Your results will depend on data volume, skew, and pool size.
Broadcast was the right call because the dimension was small and predictable; the optimizer just had to stop treating it as a large table.
Salting remains better when the dimension is too large to broadcast, such as a table with tens of millions of rows.
How do you force a broadcast join in PySpark?
Wrap the small DataFrame in broadcast() from pyspark.sql.functions, and Spark will plan a broadcast hash join regardless of the auto-broadcast threshold, as long as the join type supports it.
Here is the production transform after the change:

Why this works: the
broadcast()function marks the dimension for broadcast, and Spark join strategy hints give a broadcast hint priority even when statistics suggest the table is above the threshold. The driver builds a hash table from the 8,000 accounts and sends it to every executor, where each task streams its own slice of events through it, so the hot accounts no longer converge on one partition. Broadcast support varies across PySpark join types: in a left join, only the right side can be broadcast, which is exactly where the dimension sits.
In Spark SQL, the same hint goes inside the query once the DataFrames are registered as views. Spark also accepts BROADCASTJOIN and MAPJOIN as aliases.
The hint carries over to other join shapes. For a PySpark join on multiple columns, pass a list of keys, such as on=["account_id", "region_id"], and Spark builds the broadcast hash table on the composite key. For a PySpark join on different column names, pass an explicit condition such as events.acct_id == accounts.account_id, then drop the dimension’s key column after the join.
We also raised the auto-broadcast threshold so similar neighboring queries would broadcast without hints. spark.conf.set covers the current session; the Spark pool’s configuration covers every notebook:

When to use which: pin known joins with the hint, and use the threshold for broad coverage, set well below what driver memory can hold.
One habit keeps every PySpark broadcast join cheap: select only the columns the join needs before broadcasting. A dimension trimmed to its key and label columns takes far less memory, which often decides whether it clears the threshold.
Chasing a similar straggler? Talk to our team about Spark performance tuning.
How do you verify the fix in the physical plan?
Check the physical plan: a working PySpark broadcast join shows a BroadcastExchange feeding a BroadcastHashJoin, with the small dimension as the build side, not SortMergeJoin.

Trimmed output from our verification run (Delta tables appear as Scan parquet):

We ran three checks before promoting the change. First, the plan showed BroadcastHashJoin LeftOuter BuildRight, confirming the dimension was the broadcast side. Second, with the dimension inflated to 120 MB, unhinted neighboring queries fell back to SortMergeJoin at the 50 MB threshold, and the hinted job still ran without driver out-of-memory errors. Third, we compared row counts and data hashes for each event_date partition against the pre-change output to confirm identical results.
In the Spark UI, task durations in the stage that scans and joins the events should now be even; a lone long task points to skew elsewhere.
What changed after the fix?
Runtime dropped 83%, from 3 hours 40 minutes to 38 minutes, and the skew hot spot disappeared because the event table stopped shuffling.
Spark pool consumption for this workload fell by about 60%. The downstream Power BI refresh now finishes nearly three hours earlier. Later pipeline stages, such as those in our guide to incremental data loading on Azure Synapse with PySpark, also start on time.
Engagement note: this is a composite case study drawn from ScriptsHub Technologies Spark performance engagements. Client details are anonymized, and figures are representative of results from that work.
How to apply a PySpark broadcast join step by step
Confirm the join strategy, size the small side, force the broadcast, then prove the plan changed.
Step 1: Inspect the physical plan. Run .explain() or open the Spark UI and note whether the join is SortMergeJoin, ShuffledHashJoin, or BroadcastHashJoin.
Step 2: Measure the small side. Check its in-memory size, not its row count; it must fit comfortably in driver and executor memory.
Step 3: Check the planner’s estimate. Run explain(mode="cost") to see the size Spark assumes. For Delta tables, the estimate starts from file sizes in the Delta log, and selecting fewer columns lowers it. For Parquet or Hive tables in the metastore, refresh Spark table statistics with ANALYZE TABLE ... COMPUTE STATISTICS.
Step 4: Apply the hint. Wrap the small side with broadcast(df) and confirm it sits on the side your join type allows.
Step 5: Tune the threshold. Raise spark.sql.autoBroadcastJoinThreshold to cover your typical dimensions, following Microsoft’s guidance on Apache Spark performance in Azure Synapse for memory sizing.
Step 6: Re-check the plan and test the boundary. Confirm BroadcastHashJoin appears, then rerun with an oversized dimension to make sure the driver survives. With AQE off, a broadcast slower than spark.sql.broadcastTimeout (300 seconds by default) fails the job. With AQE on, Spark no longer applies that timeout (SPARK-36414), so watch the broadcast stage’s duration in the Spark UI.
Step 7: Monitor spill. Watch the Spark UI for memory spill and raise executor memory if it appears.
When should you use a broadcast join instead of a bigger cluster?
Whenever one side of a skewed join is small and predictable, fix the join strategy before you add compute.
A larger pool would only have finished the other 199 tasks sooner while the straggler kept running on one core. A correctly sized PySpark broadcast join, backed by accurate size estimates and verified in the physical plan, removes shuffle-driven skew without touching the schema or business logic. Our other posts on timezone conversion in Azure Synapse, late-arriving data with watermarks, and ADLS Gen2 access control cover the surrounding pipeline, and the full library of data engineering guides is on our blog.
If skewed Spark jobs are costing your team its morning window, ScriptsHub Technologies can diagnose the plan and deliver a tuned join strategy in a short, fixed-scope engagement. Start the conversation about data engineering consulting with our team.
Frequently asked questions
Q. What is broadcast join in PySpark?
A PySpark broadcast join copies the smaller DataFrame to every executor, so the larger one is joined locally with no network shuffle. It suits small lookup or dimension tables.
Q. What is the difference between shuffle join and broadcast join?
A shuffle join repartitions both tables by the join key across the network. A PySpark broadcast join sends the small table to every executor, so the large table is joined where it already sits.
Q. How does a broadcast hash join work in Spark?
The driver collects the small table, builds an in-memory hash table from it, and broadcasts that table to every executor. Each task then streams its partition of the large table through it.
Q. How to use joins in PySpark?
Use df1.join(df2, on="key", how="inner") to join two DataFrames in PySpark. The how argument sets the join type, and wrapping the smaller DataFrame in broadcast() avoids shuffling the larger one.
Q. What are the different types of joins in Spark?
Spark supports inner, left, right, full outer, left semi, left anti, and cross joins. Not all PySpark join types can broadcast: a full outer join, for example, cannot use a broadcast hash join.
Related posts:




