Sooner or later every Spark job needs a piece of reference data on every executor: a country lookup table, a set of blocked user ids, a trained model's coefficients, a dictionary for tokenising text. The naive way to get it there is to reference a driver-side variable inside a function and let Spark ship it. That works, and it quietly multiplies network traffic and driver CPU by the number of tasks. Broadcast variables exist to ship a read-only value once per executor instead of once per task.
This article explains what actually happens when you call broadcast(): how the value is serialized, cut into pieces and spread through the cluster, where the copies live in memory, how they are cleaned up, and how to refresh one in a long-running job. It ends with a worked enrichment job, the failure modes you will meet in production and a checklist. Broadcast joins in Spark SQL use the same machinery and are covered separately in the broadcast join article.
The problem broadcast solves: closures ship per task
When you pass a function to map, filter or mapPartitions, Spark serializes the function together with every variable it captures: its closure. Since Spark 1.1 the DAGScheduler serializes each stage's RDD and function once, as the stage's task binary, and distributes it with the same broadcast machinery, so a captured 200 MB dictionary is not sent once per task. It is still expensive: it is serialized and shipped again for every stage and every job that captures it, and every task deserializes its own copy, so an executor running eight tasks at once pays eight deserializations and briefly holds eight copies.
# Anti-pattern: the dict is captured by the closure and serialized into EVERY task.
geo = load_geo_table() # ~200 MB in memory
rdd.map(lambda ip: lookup(geo, ip)) # geo rides in the stage's task binary: re-shipped per stage,
# deserialized by every task
# Broadcast: serialized once, shipped at most once per executor, read through .value.
geo_b = sc.broadcast(load_geo_table())
rdd.map(lambda ip: lookup(geo_b.value, ip))Spark gives you a hint when this happens. When a stage's task binary exceeds 1000 KiB, the DAGScheduler logs a warning beginning Broadcasting large task binary with size. A separate warning, that a stage contains a task of very large size, fires when one serialized task exceeds the same limit, usually because of local data passed to parallelize. The fix for a large captured object is an explicit broadcast: serialized once, fetched at most once per executor, deserialized once there and shared by every task, stage and job that reads it.
Broadcast is the right tool when the value is read-only, needed by many tasks, and fits comfortably in executor memory. It is the wrong tool for data that changes mid-job, data larger than a few hundred megabytes, or results flowing back to the driver.
The API in full
The surface is small. SparkContext.broadcast(value) returns a Broadcast handle; tasks read the value through .value; and two methods control cleanup.
// Scala
val b: Broadcast[Map[String, Int]] = spark.sparkContext.broadcast(lookupMap)
val out = rdd.map(k => b.value.getOrElse(k, -1))
b.unpersist() // drop executor copies; re-fetched lazily if used again
b.unpersist(true) // same, but block until executors confirm removal
b.destroy() // drop driver AND executor copies; any later use throws
# PySpark
b = spark.sparkContext.broadcast(lookup_dict)
rdd.map(lambda k: b.value.get(k, -1))
b.unpersist(blocking=False)
b.destroy()The handle is what you capture in closures. It is tiny when serialized, essentially a broadcast id, so tasks stay small. Calling .value inside a task is what triggers the fetch on that executor.
The difference between the two cleanup calls matters. unpersist() deletes the cached copies on executors but keeps the driver's copy, so if a later stage reads the broadcast again, executors simply fetch it again. destroy() removes all data and metadata, including on the driver; any later use, including a lazily evaluated RDD that still references it, fails with an error saying the broadcast was used after it was destroyed.
Inside TorrentBroadcast
Spark has a single broadcast implementation, TorrentBroadcast (an older HTTP-based one was removed in Spark 2.0). Its design is borrowed from BitTorrent, and it exists to stop the driver being the bottleneck when hundreds of executors want the same bytes.
- On the driver, at broadcast time. The value is serialized with the configured serializer (Java by default, Kryo if you set
spark.serializer), compressed withspark.io.compression.codecwhenspark.broadcast.compressis true (the default), and cut into pieces ofspark.broadcast.blockSize(4 MB by default). The pieces are stored in the driver's BlockManager, and withspark.broadcast.checksum(default true) a checksum is recorded for each. The driver also keeps the original object for its own use. - On an executor, at first use. When a task calls
.value, the executor checks its local BlockManager for the already-assembled value. If it is missing, it takes a lock for that broadcast id, so concurrent tasks wait instead of fetching twice, and asks the driver which nodes hold each piece. - Peer fetch. The executor fetches pieces in random order, from the driver or from any other executor that already holds them. Each piece fetched is registered in the local BlockManager immediately, so this executor becomes a source for others.
- Assemble and cache. With all pieces present, the executor verifies checksums, decompresses, deserializes and stores the resulting object in its BlockManager at MEMORY_AND_DISK level.
Two properties fall out of this design. Fetching is lazy, so an executor that never runs a task touching the broadcast never fetches it, and an executor added later by dynamic allocation fetches it on demand from peers. And because one object is shared by all task threads on an executor, the object must be safe to read concurrently. Immutable maps and arrays are fine; a lazily initialised cache, a non-thread-safe parser or a regex matcher with internal state is not.
Where the memory goes
A broadcast costs memory in more places than people expect. On the driver: the original object, plus the serialized and compressed pieces, held until cleanup. Driver out-of-memory errors at the broadcast() call are almost always this.
On each executor JVM: the pieces and the deserialized object, both counted as storage memory in the unified memory pool described in Spark memory management. Under pressure the cached blocks can spill to local disk, since the level is MEMORY_AND_DISK, but the object a running task is reading is pinned on the heap. Deserialized JVM objects are often three to five times larger than their serialized form; a HashMap of boxed values is the classic offender. Broadcasting primitive arrays, or using Kryo with registered classes, shrinks both the wire size and the heap size.
PySpark multiplies the executor cost. In PySpark the value is pickled on the driver, and each Python worker process unpickles its own copy when it first touches .value. An executor with eight cores typically runs up to eight Python workers, so a 500 MB Python dictionary becomes up to 4 GB of Python memory per executor, outside the JVM heap and counted against executor overhead memory on YARN and Kubernetes. The remedies are compact representations (sorted arrays plus bisect instead of dicts, or NumPy arrays), fewer cores per executor, a larger spark.executor.memoryOverhead or, since Spark 2.4, an explicit spark.executor.pyspark.memory limit so the failure is clear rather than a container kill.
Worked example: IP-to-country enrichment
A daily job adds a country to a billion log events using a table of 3 million IPv4 ranges. As a Python list of tuples that table is several hundred megabytes; as two parallel lists of integers and short strings it is closer to 150 MB pickled. The job sorts the ranges once on the driver and broadcasts the two lists:
import bisect
from pyspark.sql import SparkSession, functions as F, types as T
spark = SparkSession.builder.getOrCreate()
sc = spark.sparkContext
# 3 million non-overlapping IPv4 ranges, sorted by start. Two parallel lists are far more
# compact than a list of tuples or a dict, and bisect needs only the starts.
starts, countries = load_ranges("geo_ranges.csv") # list[int], list[str]
ranges_b = sc.broadcast((starts, countries))
def ip_to_int(ip: str) -> int:
a, b, c, d = (int(x) for x in ip.split("."))
return (a << 24) | (b << 16) | (c << 8) | d
def country_partition(rows):
starts, countries = ranges_b.value # first call on this worker unpickles
for row in rows:
i = bisect.bisect_right(starts, ip_to_int(row.ip)) - 1
yield (row.event_id, countries[i] if i >= 0 else None)
events = spark.read.parquet("s3://logs/events/dt=2026-10-01/")
enriched = events.select("event_id", "ip").rdd.mapPartitions(country_partition)
enriched.toDF(["event_id", "country"]).write.parquet("s3://logs/enriched/dt=2026-10-01/")Run with 100 executors of 4 cores, the broadcast moves roughly 150 MB to each executor once, peer to peer, and deserializes it once per Python worker, instead of once per task across 4,000 or more tasks in every stage that captures it. Each executor holds up to four unpickled copies, one per Python worker, so overhead memory is sized for about 600 MB of Python objects per executor plus headroom.
Would a join be better? Spark SQL could join events to ranges with a range condition, but with no equality key it can use neither a hash join nor a sort-merge join, and falls back to a broadcast nested-loop join or a cartesian product unless you add an equality key such as an IP-prefix bucket. For an equality lookup on a DataFrame, prefer letting Spark SQL plan a broadcast hash join, which uses the same TorrentBroadcast underneath.
Lifecycle, cleanup and refresh
You rarely need to call unpersist in a batch job. Spark's ContextCleaner tracks Broadcast handles on the driver with weak references; when a handle becomes unreachable and the driver JVM garbage collects it, the cleaner removes the broadcast's blocks from the driver and all executors. Because that depends on driver GC, and a driver with a large heap may rarely collect, Spark also triggers a periodic GC on the driver every spark.cleaner.periodicGC.interval (30 minutes by default). Long-lived handles held in a field or a global never become unreachable, so they are never cleaned.
Long-running jobs, such as streaming applications and notebook sessions, therefore need explicit hygiene. A broadcast is immutable, so "updating" one means creating a new broadcast and dropping the old one. In Structured Streaming the natural place is foreachBatch, which runs on the driver once per micro-batch:
# Long-running streaming job: refresh a rules table every 10 minutes without leaking copies.
import time
state = {"b": sc.broadcast(load_rules()), "loaded": time.time()}
def process_batch(df, batch_id):
if time.time() - state["loaded"] > 600:
old = state["b"]
state["b"] = sc.broadcast(load_rules()) # new id, new pieces
state["loaded"] = time.time()
old.unpersist(blocking=False) # not destroy(): a task may still hold it
rules_b = state["b"] # capture the handle, not the dict
df.rdd.mapPartitions(lambda rows: apply_rules(rules_b.value, rows)) \
.toDF().write.mode("append").parquet(OUT)
query = stream.writeStream.foreachBatch(process_batch).start()Unpersist rather than destroy the old version: a task from a previous batch might still be reading it, and unpersist only removes cached copies, while destroy would fail any straggling reference. Capture the handle in a local variable before building the closure, so each batch's tasks see one consistent version.
Failure modes
- Driver OOM at broadcast time. The object, its serialized form and buffers coexist on the driver. Shrink it or join instead.
- Executor or container OOM on first use. Deserialized size is several times the serialized size, multiplied by Python workers in PySpark. Measure the deserialized object on a single node before shipping it.
- Expecting mutations to propagate. Changing
b.valueinside a task changes that executor's local copy only, silently and inconsistently. Treat broadcast values as frozen; use accumulators or a write path for results. - Capturing the enclosing object. In Scala, referencing a broadcast through a field of a class (
this.lookup.value) captures the whole outer object, which may be unserializable or huge. Copy the handle to a localvalbefore the closure. - Use after destroy. A cached RDD or DataFrame recomputed after its broadcast was destroyed fails at recompute time, possibly hours later. Destroy only at the very end, or unpersist instead.
- Non-thread-safe values. Concurrent tasks share the object. Mutable caches, date formatters and some client libraries corrupt state or crash under concurrent reads.
- Stale data in long-running jobs. A broadcast is a snapshot taken at creation. If the source changes, nothing notices until you rebroadcast.
Trade-offs: broadcast, join or lookup service
| Option | Good when | Costs |
|---|---|---|
| Broadcast variable | Read-only reference data up to a few hundred MB, custom lookup logic (ranges, tries, models) | Driver and executor memory; immutable snapshot; manual refresh |
| Broadcast hash join (SQL) | Equality join with a small side, DataFrame code | Same memory cost; planner decides by size estimates |
| Shuffle join | Both sides large, or the lookup table too big to replicate | A shuffle of both sides; skew risk |
| External lookup service or cache | Data changes continuously or is far too large | Network call per row or batch; availability coupling |
A sensible default: broadcast when the serialized value is under about 100 MB and most tasks need it. Persisting the large side is complementary: broadcast the small side, cache the big one if it is reused.
What to do next
- Search your driver logs for the large task binary warning; each hit usually points at a captured object that should be an explicit broadcast.
- For every existing broadcast, record serialized size, deserialized size on one executor, and, in PySpark, the number of Python workers per executor; multiply before you blame the cluster.
- Replace dicts of boxed values with sorted primitive arrays or NumPy arrays where lookups allow it, and register classes with Kryo for JVM jobs.
- Capture broadcast handles in local variables before building closures, never through fields of a larger object.
- In long-running and streaming jobs, rebroadcast on a schedule inside
foreachBatchand unpersist the previous version. - For equality lookups on DataFrames, let Spark SQL plan a broadcast join instead of hand-rolling one, and check the plan with EXPLAIN.