Every Spark beginner writes a counter that stays at zero. They increment a variable inside a map or foreach, print it on the driver and get nothing. The variable was serialized into each task, incremented on an executor and thrown away. Accumulators are Spark's answer: a shared variable that tasks can only add to, whose per-task updates are sent back and merged on the driver.

The trouble is that accumulators look simpler than they are. Their value depends on lazy evaluation, task retries, speculative execution and stage recomputation, and the official guarantee is narrower than most people assume. This page explains the mechanism from first principles, shows exactly when counts are right and when they drift, and covers custom accumulators, PySpark and the DataFrame-native alternative. The guarantee wording is quoted from the Spark 4.2 RDD programming guide.

Advertisement

The closure problem accumulators solve

When you pass a function to an RDD operation, Spark serializes the function and every variable it captures, and ships a copy to each task. Updates to those copies stay on the executor.

counter = 0
def count_bad(line):
    global counter
    if not line.startswith("{"):
        counter += 1          # increments a copy inside the Python worker

rdd = sc.textFile("s3://logs/2026-10-01/*.json")
rdd.foreach(count_bad)
print(counter)                # 0 on the driver, whatever the data

bad = sc.accumulator(0)
rdd.foreach(lambda line: bad.add(1) if not line.startswith("{") else None)
print(bad.value)              # the real count

The Scala equivalent is worse: in local mode a JVM closure can appear to work because the driver and executor share one JVM, so the bug survives until the job reaches a cluster. Accumulators exist to make the update path explicit: tasks get a zero-valued local copy, add to it, and the driver merges the copies.

The API

In Scala and Java, SparkContext creates the built-in kinds: sc.longAccumulator("name"), doubleAccumulator and collectionAccumulator. LongAccumulator also tracks a count and exposes sum, count and avg. Custom types extend AccumulatorV2[IN, OUT] and are registered with sc.register(acc, "name"). In Python, sc.accumulator(0) covers numbers and custom types subclass AccumulatorParam.

Name your accumulators. Named accumulators appear in the web UI on the stage that modifies them, with each task's contribution in the Tasks table, which turns a mysterious count into something you can trace to a partition; the Spark UI guide shows where.

Advertisement

How an update travels

One accumulator, one job: copies go out with tasks, updates come back with resultsDriverregister acc (id 7, value 0)DAGScheduler merges updatesof successful tasks onlyTask 0 (executor A)local copy starts at zeroTask 1 (executor B)local copy starts at zeroTask 1, attempt 2retry after executor lossserializeadd() x 120result: +120add() x 40, then lostupdate discardedadd() x 95result: +95Task results carry accumulator updates back to the driverdriver: 0 + 120 + 95 = 215Inside an action: each partition's update is applied once.Inside a transformation: a recomputed stage adds its updates again.The driver only ever sees whole-task updates. Reading the value on an executor shows a partial, not the total.
Copies leave with tasks; only successful task results are merged on the driver.

The lifecycle has four steps. The driver registers the accumulator, giving it an ID. When a task is serialized, the accumulator goes with it, and on deserialization the executor creates a fresh copy with the zero value. As the task runs, add changes only that copy. When the task finishes, its accumulator updates travel back with the task result, and the driver's scheduler merges them into the registered instance when it handles the task completion event.

Two details follow from this, and both come from Spark internals rather than the user guide. Updates from failed tasks are discarded for user accumulators, so a task that dies halfway contributes nothing. And executors also send in-progress values with heartbeats so the UI can show live metrics, but those do not change the value your driver code reads. The upshot is that acc.value on the driver reflects completed tasks only, and is reliable only after the action that ran them has returned. The scheduling units involved are explained in Spark stages and tasks.

The guarantee, and what breaks it

The programming guide states it precisely: "For accumulator updates performed inside actions only, Spark guarantees that each task's update to the accumulator will only be applied once, i.e. restarted tasks will not update the value. In transformations, users should be aware of that each task's update may be applied more than once if tasks or job stages are re-executed." foreach and foreachPartition are actions; map, filter and flatMap, and UDFs in DataFrame plans, are transformations.

SituationUpdate inside an actionUpdate inside a transformation
No action is ever runNothing runsStays at zero: transformations are lazy
Task fails and is retriedCounted onceFailed attempt discarded; usually once
Speculative duplicate also finishesCounted onceCan be counted twice
Lost shuffle output forces a parent stage to rerunCounted onceParent stage updates added again
Two actions on an uncached RDDEach action counts its own runCounted once per action
Cached RDD partially evicted, then reusedNot applicableEvicted partitions recomputed and counted again

The last two rows surprise people most. If you count bad records in a map, then call count() and later saveAsTextFile(), the map runs twice and the counter doubles. Caching reduces this but does not remove it, because cache eviction or executor loss forces recomputation.

Worked example: counting bad records

A job parses 10 million log lines in 200 partitions, of which 12,400 are malformed. Version one counts inside a transformation:

bad = sc.accumulator(0)

def parse(line):
    try:
        return [json.loads(line)]
    except ValueError:
        bad.add(1)
        return []

events = sc.textFile(path, 200).flatMap(parse)
n = events.count()                          # bad.value == 12400
events.map(to_row).saveAsTextFile(out)      # bad.value == 24800: flatMap ran again

Caching events before count() brings the final figure back to 12,400, until an executor holding 3 cached partitions is lost during the save. Those 3 partitions are recomputed from the source and their bad lines, roughly 186 at an even spread, are added again, giving about 12,586. The job succeeded and the output is correct; only the metric is wrong.

Version two moves the count into an action over the same data, so retries and recomputation cannot inflate it:

bad = sc.accumulator(0)
parsed = sc.textFile(path, 200).map(lambda l: (try_parse(l), l)).cache()

parsed.foreachPartition(lambda it: bad.add(sum(1 for rec, _ in it if rec is None)))
# exactly 12400, once per partition, even with retries
parsed.filter(lambda x: x[0] is not None).map(lambda x: to_row(x[0])).saveAsTextFile(out)

Version three, often the best, does not use an accumulator: write malformed lines to a quarantine path and count them with an ordinary aggregation. That count is part of the data, is exact, and leaves evidence you can inspect. Spark debugging covers isolating bad records in this way.

Writing a custom AccumulatorV2

A custom accumulator implements six methods: isZero, copy, reset, add, merge and value. The merge must be commutative and associative, because tasks finish in any order, and the state should stay small, because every task ships its copy back to the driver. A useful example is a latency histogram with fixed buckets, which is small and merges by adding arrays.

import org.apache.spark.util.AccumulatorV2

class HistogramAcc(bounds: Array[Double]) extends AccumulatorV2[Double, Array[Long]] {
  private val counts = new Array[Long](bounds.length + 1)

  override def isZero: Boolean = counts.forall(_ == 0L)
  override def copy(): HistogramAcc = {
    val c = new HistogramAcc(bounds)
    System.arraycopy(counts, 0, c.counts, 0, counts.length)
    c
  }
  override def reset(): Unit = java.util.Arrays.fill(counts, 0L)
  override def add(v: Double): Unit = {
    val i = java.util.Arrays.binarySearch(bounds, v)
    counts(if (i >= 0) i else -i - 1) += 1
  }
  override def merge(other: AccumulatorV2[Double, Array[Long]]): Unit = other match {
    case h: HistogramAcc => for (i <- counts.indices) counts(i) += h.counts(i)
    case _ => throw new UnsupportedOperationException(other.getClass.getName)
  }
  override def value: Array[Long] = counts.clone()
}

val latency = new HistogramAcc(Array(10.0, 50.0, 100.0, 500.0))
sc.register(latency, "request_latency_ms")
requests.foreach(r => latency.add(r.latencyMs))   // action: counted once

isZero must be true after reset, because Spark creates each task's copy with copyAndReset and checks that the result is zero. Avoid collectionAccumulator and list-shaped custom types for anything that grows with the data: every element is shipped to and held in driver memory.

PySpark specifics

In Python, each worker process adds to its local copy and the updates come back with the task result to the driver's Python process. Custom types subclass AccumulatorParam with zero(value) and addInPlace(v1, v2). Reading .value inside a task raises an error rather than returning a partial, which is helpful. The same action and transformation rules apply, and Python UDFs in DataFrame plans count as transformations.

Spark Connect clients have no SparkContext, so the RDD accumulator API is not available there. Code that may run through Spark Connect, including many managed shared clusters, should use the DataFrame approach below.

The DataFrame alternative: observe

For DataFrame jobs, DataFrame.observe with an Observation (Spark 3.3 and later in both Scala and Python) computes named aggregate expressions as a side effect of the next action on that plan, and returns them on the driver. The metrics belong to one query execution, so a second action does not add to them, which removes the commonest recount. In PERMISSIVE mode from_json turns a malformed line into a struct of null fields rather than a null struct, so the example tests a field every valid record has.

from pyspark.sql import Observation
from pyspark.sql import functions as F

obs = Observation("ingest_quality")
parsed = (spark.read.text(path)
          .withColumn("j", F.from_json("value", schema))
          .observe(obs,
                   F.count(F.lit(1)).alias("rows"),
                   F.count(F.when(F.col("j.ts").isNull(), 1)).alias("bad")))

parsed.where("j.ts IS NOT NULL").select("j.*").write.parquet(out)
print(obs.get)        # {'rows': 10000000, 'bad': 12400}

Observations are still collected through accumulator machinery inside the plan, and nothing in the documentation promises they are immune to stage re-execution, so when a number must be exact, keep the quarantine-and-aggregate approach. Use observations for data quality metrics and accumulators for things aggregates cannot express cheaply, such as timing inside a function. For broadcast, the other shared variable, see Spark broadcast variables.

Operating accumulators in real jobs

Spark itself is the heaviest user of accumulators. Task metrics such as bytes read and shuffle records are implemented on the same AccumulatorV2 machinery, and so are the SQL metrics shown on plan nodes in the SQL tab. Internal metrics are marked to count values from failed tasks too, because for diagnosis you want to see work that was wasted; user accumulators do not, which is the behaviour described above. That shared plumbing means accumulators are cheap enough for a handful of counters per job, but every registered accumulator touched by a task travels back with every task result, so a job with 100,000 tasks and dozens of large custom accumulators puts real load on the driver.

For metrics that leave the job, read the value on the driver after the action returns and publish it to your monitoring system with the job and run identifiers. Do not publish from inside tasks: the task may be retried and the metric duplicated in a system that cannot tell. In Structured Streaming, foreachBatch gives you a driver-side function per micro-batch; create the counter outside, call reset() on a custom accumulator or record the difference between successive values, and publish once per batch ID. A restarted query can replay a batch, so key published metrics by batch ID and let the monitoring side overwrite rather than add.

Finally, treat accumulator totals as observability, not as data. If a number appears on an invoice, in an audit or in a downstream decision, compute it as part of the dataset, where Spark's fault tolerance guarantees it.

Failure modes

  • Counter stays at zero: no action ran, or a plain variable was used instead of an accumulator.
  • Counter doubles: updates in a transformation and two actions, or recomputation after cache loss.
  • Driver out of memory: a collection accumulator gathering per-row values from billions of rows.
  • Non-deterministic custom totals: a merge that depends on order, such as keeping the first value seen.
  • Value read too early: reading on the driver from another thread while the action is still running gives a partial count.
  • Metrics used for business logic: a billing or audit number taken from an accumulator in a transformation is wrong whenever any task reruns.

What to do next

  1. Search your jobs for accumulator updates inside map, filter, flatMap and UDFs, and list which ones feed alerts or reports.
  2. Move those that must be exact into foreach or foreachPartition or an ordinary aggregation, and use DataFrame.observe for monitoring metrics.
  3. Name every accumulator so its per-task values show in the UI.
  4. Replace collection accumulators with bounded types such as counters or fixed-bucket histograms.
  5. Test the count under failure: kill an executor mid-job in staging and compare the accumulator with the true count.
  6. Check whether any job runs through Spark Connect, and migrate its accumulators to observations first.
Key takeaway: Accumulators exist because closures copy variables into tasks; each task adds to a zero-valued copy and the driver merges the copies from successful tasks. Spark guarantees each task's update is applied once only for updates made inside actions; in transformations, retries, speculation, recomputation and repeated actions can count again. Update inside foreach or foreachPartition when the number must be exact, keep custom accumulators small and commutative, and prefer DataFrame.observe or a real aggregation for data quality metrics.