What you will be able to do
- Explain why a normal driver-side variable used inside a Spark function is copied to each task, and why changes made there never reach the driver
- Describe broadcast variables (read-only, cached on each node), and how to create, read, unpersist and destroy them
- Describe accumulators (add-only, readable only on the driver), and explain why updates inside lazy transformations are unreliable
- Tell a broadcast variable apart from the broadcast() join hint and from SQL session variables, and know where SparkContext-based shared variables are not available on Databricks
Key concept
Shared variables — By default, each task gets its own private copy of any driver variable it uses. Spark gives you two limited ways around that. A broadcast variable sends a read-only value to every node once. An accumulator lets tasks add to a value that only the driver reads.
1.Why Spark needs special variables at all
A Spark application has two kinds of process: a driver that runs your main program, and executors on worker nodes that run tasks. When you pass a function to a Spark operation such as map, filter or foreach, Spark serializes the function together with every variable it refers to (its closure) and sends that bundle to the executors. Each task works on its own copy, so a variable used this way acts like a local value in each task. Reading it works. Writing to it changes only that task's copy. The Spark programming guide puts it plainly: no updates on the remote machine are propagated back to the driver program.
This default causes two problems that come up all the time. First, a large lookup structure such as a dictionary of country codes or a reference list is shipped again and again with tasks, which wastes network and memory. Second, you sometimes want tasks to report something back to the driver, such as how many malformed records they skipped, and an ordinary variable cannot do that. General read-write variables shared across all tasks would be too expensive to coordinate, so Spark offers two narrow tools instead. Broadcast variables fix the first problem. Accumulators fix the second. Together they are called shared variables, and they are the two variable types this exam objective is about.
Checkpoint 1 of 8· Check yourself
Which statement best describes a SQL session variable created with DECLARE VARIABLE on Databricks?
Session variables are a SQL feature scoped to a session. They are not among the two shared-variable types (broadcast variables and accumulators) that Spark uses to share data between the driver and tasks.
“Variables are typed and schema qualified objects which store values that are private to a session.”Source: docs.databricks.com
2.Broadcast variables: one read-only copy per node
A broadcast variable wraps a value that you want every executor to read. Spark sends it to each node once and keeps it cached there, so tasks no longer carry it with them. You create one from the SparkContext (in PySpark, spark.sparkContext.broadcast(v)), and both the driver and tasks read it through its .value attribute. The value is meant to be read-only. Tasks should not change it, and if the driver changes its local Python object after broadcasting, executors keep the copy they already have.
The Spark guide adds a nuance. Spark already broadcasts the common data that tasks within one stage need. So creating a broadcast variable yourself pays off mainly when the same data is reused by tasks across several stages, or when you want it cached in deserialized form. A typical case is a medium-sized lookup table used inside a UDF or an RDD map, which would otherwise be serialized into every task's closure.
b = spark.sparkContext.broadcast([1, 2, 3, 4, 5])Checkpoint 2 of 8· Fill the gap
Which SparkContext method completes this line so that b is a read-only value cached on every executor?
b = spark.sparkContext. ? ([1, 2, 3, 4, 5])SparkContext.broadcast() returns a pyspark.Broadcast object, and you read its contents through .value. accumulator() creates an add-only variable, and parallelize() turns a local list into an RDD.
Source: spark.apache.orgBroadcast variables have a lifecycle. unpersist() drops the cached copies on the executors. If you use the variable again afterwards, Spark broadcasts it again. destroy() removes all of its data and metadata for good, and the variable can't be used after that. By default neither call blocks; pass blocking=True if you need to wait until the resources are freed.
One naming collision catches many candidates. pyspark.sql.functions.broadcast(df) is not a broadcast variable. It is a join hint that marks a small DataFrame so the optimizer can use a broadcast join. Both send data to every executor, but one is a DataFrame query hint and the other is a SparkContext-level value that you read in your own code.
Checkpoint 3 of 8· Match them up
Match each API to what it does
Tap a term, then the definition that fits it.
The first three belong to the broadcast variable lifecycle on pyspark.Broadcast. The SQL function broadcast(df) is a join hint on a DataFrame, which is a separate feature.
“Marks a DataFrame as small enough for use in broadcast joins.”Source: docs.databricks.com
Checkpoint 4 of 8· Exam question
A data engineer has a small Python dictionary `country_codes` (well under 1 MB) that maps ISO codes to country names, and needs to look it up inside a PySpark UDF applied to a 500 million row DataFrame. To avoid Spark serializing and shipping a fresh copy of `country_codes` with every task, which change should the engineer make to the following code? ```python def lookup_country(code): return country_codes.get(code, "Unknown") lookup_udf = udf(lookup_country, StringType()) df = df.withColumn("country_name", lookup_udf(col("iso_code"))) ```
Correct answer: A — Wrap `country_codes` with `spark.sparkContext.broadcast(country_codes)` and reference `bc.value.get(code, "Unknown")`, so every executor keeps one cached copy.
- A. Broadcasting caches a single read-only copy of the dictionary on each executor, and referencing `.value` inside the UDF avoids re-serializing and re-sending the dictionary with every task.
- B. Rebuilding the dictionary inside the function makes the problem worse: Spark still ships whatever data built it, and the work of constructing the mapping now repeats on every single row.
- C. A join is a valid pattern for larger reference tables, but for a small in-memory dictionary it adds an unnecessary shuffle and does not fix the serialization behavior of the existing UDF code.
- D. Accumulators are write-only aggregators whose running total only the driver can read; they are not a mechanism for distributing an arbitrary read-only mapping for executors to look values up in.
- E. `spark.sql.autoBroadcastJoinThreshold` only influences whether Catalyst chooses a broadcast join strategy for DataFrame joins, and has no effect on a plain Python closure referenced inside a UDF.
3.Accumulators: tasks add, only the driver reads
An accumulator solves the opposite problem: it carries information from tasks back to the driver. You create one with sc.accumulator(initial_value). Inside tasks you can only add to it, with .add(x) or +=. Tasks cannot read it; only the driver reads the merged total through .value. Updates are combined with an associative and commutative operation, so Spark can merge partial results from many tasks in any order. That is why accumulators fit counters and sums. Spark supports numeric accumulators out of the box, and you can define custom types (in Scala and Java, by subclassing AccumulatorV2). If you give an accumulator a name, the Spark UI shows it for the stage that updates it, but the guide notes this is not yet supported in Python.
accum = sc.accumulator(0)
def g(x):
accum.add(x)
return f(x)
data.map(g)
# Here, accum is still 0 because no actions have caused the `map` to be computed.Laziness is only half the story. Even after an action runs, where you update an accumulator decides how far you can trust it. When the update happens inside an action (for example foreach), Spark guarantees each task's update is applied only once, so restarted tasks do not count twice. When the update happens inside a transformation, a retried task, a re-executed stage or a recomputed uncached RDD can apply the same update more than once. A count taken in map can come out too high. Also, if merging a task's accumulator update fails, Spark ignores the failure and still marks the task successful. So a job can succeed while an accumulator is wrong. Treat accumulators as diagnostics and counters, not as a source of business results.
| Aspect | Broadcast variable | Accumulator |
|---|---|---|
| Created with | SparkContext.broadcast(v) | SparkContext.accumulator(v) |
| Direction of data | Driver to every node | Tasks to driver |
| What tasks may do | Read via value | Only add (add or +=); cannot read |
| Who reads the result | Driver and tasks, via value | Driver only, via value |
| Cleanup | unpersist, destroy | None needed for the exam |
Checkpoint 5 of 8· Check yourself
A task running on an executor tries to read an accumulator's current total to decide whether to skip a record. What does the Spark programming guide say?
Accumulators are write-only (add-only) from a task's point of view. Only the driver program reads the merged value, so per-record logic must not depend on it.
“Tasks running on a cluster can then add to it using the add method or the += operator. However, they cannot read its value.”Source: spark.apache.org
Checkpoint 6 of 8· Exam question
A team runs the following PySpark job to count rows where `amount` is negative, then submits it as a batch job: ```python neg_count = spark.sparkContext.accumulator(0) def flag_negative(row): if row.amount < 0: neg_count.add(1) orders_df.foreach(flag_negative) print(neg_count.value) ``` The job runs once with no task failures or retries. What does `print(neg_count.value)` output?
Correct answer: A — The exact count of rows with a negative `amount`, because `foreach` is an action, and Spark guarantees exactly-once accumulator updates here.
- A. `foreach` is an action, so it forces evaluation, and without retries each task's accumulator update is guaranteed to apply exactly once, giving an exact negative-row count.
- B. Once an action like `foreach` completes, the driver can immediately read the accumulator's `.value`; a second unrelated action is not required to make the update visible.
- C. Spark does not re-execute a successfully completed task twice by default; double execution only happens under retries or speculative execution, neither of which occurred here.
- D. PySpark's default numeric accumulator created with `spark.sparkContext.accumulator(0)` is registered automatically and its `.value` is readable on the driver without any extra registration call.
- E. The driver can read an accumulator's `.value` once the action that updates it has finished; it is not blocked simply because other unrelated jobs are running in the cluster.
4.Where shared variables are available on Databricks
Both shared-variable types are created through the SparkContext (sc or spark.sparkContext), which belongs to classic Spark, alongside the RDD API. That matters on Databricks. Serverless compute supports only Spark Connect APIs. In Spark Connect the client talks to a remote Spark server and has no SparkContext, so sc.broadcast and sc.accumulator are not available there. Standard (shared) access mode compute also limits SparkContext and RDD use. If you need broadcast variables or accumulators, run on dedicated (classic) compute. Otherwise, use DataFrame-native options such as the broadcast() join hint or a plain aggregation to count rows.
Checkpoint 7 of 8· Check yourself
A notebook that calls spark.sparkContext.accumulator(0) runs fine on a dedicated cluster but fails on Databricks serverless compute. What is the most likely reason?
Accumulators and broadcast variables come from the SparkContext. Serverless runs on Spark Connect, which does not expose the SparkContext or the RDD API.
“Only Spark Connect APIs are supported. Spark RDD APIs are not supported.”Source: docs.databricks.com
Checkpoint 8 of 8· Exam question
What is the primary motivation for using a broadcast variable instead of relying on Spark's default task closure mechanism to make a value available to executors?
Correct answer: A — It caches a read-only copy of the value once per executor, avoiding the cost of resending that same value with every task needing it.
- A. This is the documented purpose of broadcast variables: keep a read-only copy cached on each machine rather than shipping a fresh copy of the value with every task.
- B. Broadcast variables are explicitly read-only after creation; changes made on an executor are local and are never synchronized back to the driver or to other executors.
- C. Broadcast variables live in executor memory for the life of the Spark application and have no relationship to the Delta transaction log or to durability across separate application runs.
- D. Broadcasting keeps the full value intact on every executor rather than splitting it into partitions; partitioning a value across the cluster describes a distributed Dataset, not a broadcast variable.
- E. Broadcast variables are held in executor memory (and can spill to local disk under pressure), and creating one has no requirement to first write the value to durable cloud storage.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.A broadcast variable is a shared mutable value: if a task updates it, other tasks and the driver see the change.Why is that wrong?
A broadcast variable is a read-only copy cached on each machine. Tasks only read it through .value, and nothing they do is sent back.
2.pyspark.sql.functions.broadcast(df) creates a broadcast variable you can read with .value.Why is that wrong?
functions.broadcast(df) is a join hint. It returns a DataFrame marked for a broadcast join. Broadcast variables come from SparkContext.broadcast().
3.An accumulator incremented inside map() always gives an exact count once the job finishes.Why is that wrong?
Updates in transformations run only when an action computes them, and they can be applied more than once if tasks or stages are retried. Spark guarantees exactly-once updates only inside actions.
4.Tasks can read an accumulator's running total to make decisions.Why is that wrong?
Tasks can only add to an accumulator. Only the driver can read its value.
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“it ships a copy of each variable used in the function to each task.”
↩︎ Why Spark needs special variables at all“no updates to the variables on the remote machine are propagated back to the driver program.”
↩︎ Why Spark needs special variables at all“explicitly creating broadcast variables is only useful when tasks across multiple stages need the same data”
↩︎ Broadcast variables: one read-only copy per node“They can be used to implement counters (as in MapReduce) or sums.”
↩︎ Accumulators: tasks add, only the driver reads“Accumulators do not change the lazy evaluation model of Spark.”
↩︎ Accumulators: tasks add, only the driver reads“restarted tasks will not update the value.”
↩︎ Accumulators: tasks add, only the driver reads“a buggy accumulator will not impact a Spark job, but it may not get updated correctly although a Spark job is successful.”
↩︎ Accumulators: tasks add, only the driver reads“Spark supports two types of shared variables: broadcast variables, which can be used to cache a value in memory on all nodes, and accumulators”
↩︎ Key concept“Broadcast variables allow the programmer to keep a read-only variable cached on each machine rather than shipping a copy of it with tasks.”
↩︎ Exam trap 1“if tasks or job stages are re-executed.”
↩︎ Exam trap 3“Tasks running on a cluster can then add to it using the add method or the += operator. However, they cannot read its value.”
↩︎ Exam trap 4“accumulator updates are not guaranteed to be executed when made within a lazy transformation like map().”
↩︎ Prediction“Tasks running on a cluster can then add to it using the add method or the += operator. However, they cannot read its value.”
↩︎ Checkpoint - 2.
“Variables are typed and schema qualified objects which store values that are private to a session.”
↩︎ Why Spark needs special variables at all - 3.https://spark.apache.org/docs/latest/api/python/reference/api/pyspark.Broadcast.htmlSecondary source
“A broadcast variable created with SparkContext.broadcast(). Access its value through value.”
↩︎ Broadcast variables: one read-only copy per node“Delete cached copies of this broadcast on the executors.”
↩︎ Broadcast variables: one read-only copy per node - 4.
“Marks a DataFrame as small enough for use in broadcast joins.”
↩︎ Broadcast variables: one read-only copy per node“Marks a DataFrame as small enough for use in broadcast joins.”
↩︎ Exam trap 2 - 5.
“Only Spark Connect APIs are supported. Spark RDD APIs are not supported.”
↩︎ Accumulators: tasks add, only the driver reads“Only Spark Connect APIs are supported. Spark RDD APIs are not supported.”
↩︎ Where shared variables are available on Databricks - 6.
“Spark Context (sc), spark.sparkContext, and sqlContext are not supported for Scala”
↩︎ Where shared variables are available on Databricks“RDD APIs are not supported.”
↩︎ Where shared variables are available on Databricks