Broadcast variables and accumulators
Why driver variables never change inside tasks, and the two shared variables Spark offers instead.
On this page
The closure problem
When a tasktask: The work for one partition in one stage, run on one CPU core. Learn more → runs a Python UDFUDF: User-defined function: your own Python function applied to each row. Flexible but much slower than built-in functions. Learn more → or an RDD function, Spark serializes the function together with every outside variable it references (its closure) and sends that copy with the task. Each executorexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more → changes its own copy; the driverdriver: The process that runs your program, plans the work and sends tasks to executors. Learn more →'s variable never changes.
bad_rows = 0 def parse(line): global bad_rows if not line.strip(): bad_rows += 1 # changes a copy on the executor return line.split(",") rdd.map(parse).count() print(bad_rows) # still 0
The same goes for appending to a driver-side list or dict inside a UDF. The second cost of closures is size: a 50 MB lookup dict referenced in a UDF is shipped with every task. With 2,000 tasks that is 100 GB sent over the network.
Broadcast variables
rates = {"EUR": 1.08, "GBP": 1.27, "INR": 0.012} # built on the driver b_rates = spark.sparkContext.broadcast(rates) @F.udf("double") def to_usd(amount, currency): return amount * b_rates.value.get(currency, 1.0) # read on executors orders.withColumn("amount_usd", to_usd("amount", "currency")) b_rates.unpersist() # drop executor copies (re-sent if used again) b_rates.destroy() # remove everywhere; unusable afterwards
- The value is serialized once and fetched by each executor the first time a task needs it, then kept in the executor's block manager for all later tasks. Executors fetch pieces from the driver and from each other (a BitTorrent-like protocol), so the driver is not the only source.
- It is read-only. Changing
.valueon an executor changes only that executor's copy. - It counts against storage memory on every executor, and must first fit on the driver.
| Broadcast variable | Broadcast join | |
|---|---|---|
| What is shipped | Any Python or JVM object | A small DataFrame |
| Who uses it | Your UDF or RDD code, via .value | Spark's join operator |
| How you get one | sparkContext.broadcast(obj) | Automatic under the threshold, or F.broadcast(df) |
| Prefer when | Code that cannot be expressed as a join (a parser, a model, a regex table) | Anything that can be a join: it stays inside CatalystCatalyst: Spark's query optimizer. It rewrites your DataFrame code into a faster equivalent plan before running it. Learn more → |
Accumulators
An accumulator is a variable tasks can only add to and only the driver can read. Spark merges every task's updates and sends them to the driver with task results.
bad_rows = spark.sparkContext.accumulator(0) def check(row): if row.amount is None: bad_rows.add(1) orders.foreach(check) # an action: updates counted exactly once print(bad_rows.value)
When accumulators over-count
Inside actions
foreach, foreachPartition: Spark applies each task's update only once, even if the task is retried.Inside transformations
map, UDFs, filter: updates can be applied again when a task is retried, a speculative copy runs, a lost stagestage: A group of steps Spark can run without moving data between machines. A new stage starts at every shuffle. Learn more → is recomputed, or the DataFrame is computed by a second action. Treat these values as approximate.Accumulators are also lazy: their value stays at zero until an action runs the transformation that updates them. In Python, custom types (sets, lists, dicts) need an AccumulatorParam with zero and addInPlace.
In DataFrame code, prefer these
- Counting with an aggregation:
F.sum(F.when(cond, 1).otherwise(0))is exact, retry-safe and optimised by Catalyst. - Observation (PySpark 3.3+): collect metrics during the action you were going to run anyway, with no extra job.
- A broadcast joinbroadcast join: A join where the small table is copied to every machine, so the big table never has to be shuffled. Learn more → instead of a broadcast variable whenever the lookup is a table.
from pyspark.sql import Observation obs = Observation("quality") checked = orders.observe(obs, F.count(F.lit(1)).alias("rows"), F.sum(F.col("amount").isNull().cast("int")).alias("null_amounts")) checked.write.mode("overwrite").parquet("out") print(obs.get) # {'rows': ..., 'null_amounts': ...}
SparkContext, so sparkContext.broadcast and sparkContext.accumulator are not available. Use a broadcast join, observe, or an aggregation instead.Common mistakes
- Incrementing a Python global inside a UDF — Each executor updates its own copy; the driver sees nothing.
- Using accumulator values from transformations as exact numbers — Retries, speculation and recomputation can count rows twice.
- Referencing a big dict directly in a UDF — It is serialized into every task. Broadcast it, or make it a broadcast join.
- Reading an accumulator before any action — Transformations are lazy, so it is still zero.
What you learned
- Why updating a driver variable inside a task does nothing
- How broadcast variables ship read-only data once per executor
- How accumulators send counters back to the driver, and when they over-count
- What to use instead in DataFrame code and Spark Connect
Key takeaways
- Closures are copied to executors; changes never return to the driver.
- Broadcast variables send read-only data once per executor instead of once per task.
- Accumulators are add-only on executors and readable only on the driver.
- Only accumulator updates in actions are guaranteed exactly once.
- In DataFrame code, prefer broadcast joins, aggregations and observe().
Check yourself
3 questionsA UDF increments a driver global. What does the driver see after the action?
Show the answer
The original value. Executors update serialized copies of the variable.
Where are accumulator updates guaranteed to be applied exactly once?
Show the answer
In actions such as foreach. For actions Spark applies each task's update once; transformations may be re-executed.
How often is a broadcast variable sent to an executor?
Show the answer
Once per executor. It is fetched on first use and cached in the executor's block manager.
Keep going
Up next · lesson 11 of 30 · 3 min readShuffle partitions
The default of 200 shuffle partitions, and how to size them for your data.
Related lessons
Broadcast joinsPySpark functions · 4 min read
UDFs and pandas UDFsSpark internals · 3 min read
Spark Connect
Previous: Caching and persistence
Primary sources: RDD guide: shared variables · RDD guide: understanding closures