Skip to content
Good engineers know 5 min read

Broadcast variables and accumulators

Why driver variables never change inside tasks, and the two shared variables Spark offers instead.

TL;DR Code that runs on executors gets its own copy of every driver variable it references, so changes never come back. Spark has two shared variables for this: broadcast variables send read-only data (a lookup dict, a model) to each executor once, and accumulators let tasks add to a value that only the driver can read. Accumulator updates made inside transformations can be applied more than once.

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.

PySpark · A counter that stays at zero
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

PySpark · Ship a lookup table once per executor
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 .value on an executor changes only that executor's copy.
  • It counts against storage memory on every executor, and must first fit on the driver.
Broadcast variableBroadcast join
What is shippedAny Python or JVM objectA small DataFrame
Who uses itYour UDF or RDD code, via .valueSpark's join operator
How you get onesparkContext.broadcast(obj)Automatic under the threshold, or F.broadcast(df)
Prefer whenCode 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.

PySpark · Counting bad rows correctly
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.
PySpark · Metrics without an extra job
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': ...}
Note: Spark Connect clients have no 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 questions

A 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 read
Shuffle partitions
The default of 200 shuffle partitions, and how to size them for your data.

Related lessons

Previous: Caching and persistence

Primary sources: RDD guide: shared variables · RDD guide: understanding closures