Skip to content

[SPARK-59845][CORE] Register an accumulator with the TaskContext only after it is fully deserialized - #59124

Closed
liuzqt wants to merge 2 commits into
apache:masterfrom
liuzqt:accumulator-register-after-deserialization
Closed

liuzqt wants to merge 2 commits into
apache:masterfrom
liuzqt:accumulator-register-after-deserialization

Conversation

@liuzqt

@liuzqt liuzqt commented Sep 29, 2026

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

Move the TaskContext registration out of readObject and into readResolve.

AccumulatorV2 currently registers from readObject, which runs at AccumulatorV2's own slot. Java serialization reads hierarchy data superclass-first, so no subclass field is assigned yet — while registerAccumulator publishes this into TaskMetrics where other threads can read it. The PR drops readObject (beyond registering, all it did was in.defaultReadObject(), which the framework does anyway) and registers from final protected def readResolve() instead, which runs only after every class in the hierarchy has been read.

It's final for the same reason writeReplace already is: readObject must be private so it's invoked per class, whereas readResolve is inherited — a subclass defining one would silently stop registration.

Why are the changes needed?

Subclasses are published before their state exists, so a subclass holding state in a val exposes a null. The reader that hits it is the executor heartbeater, calling isZero on every registered accumulator. The NPE propagates out of a scheduleAtFixedRate task, so ScheduledThreadPoolExecutor stops rescheduling — the executor never heartbeats again, the driver removes it on timeout.

Plus the JLS 17.5.3 aspect: a val is a final field written reflectively after publication, so a reader has no guarantee of seeing it and the JIT may fold the read.

SPARK-20977 fixed this locally in CollectionAccumulator; every subclass since has re-derived the workaround and not all have — SetAccumulator in sql/execution/debug still dereferences private val _set from isZero, and EventTimeStatsAccum escapes the NPE only incidentally via Scala's null-safe ==, reporting the wrong answer instead.

Does this PR introduce any user-facing change?

NO

How was this patch tested?

New test case

Was this patch authored or co-authored using generative AI tooling?

Yes. Claude Opus 5.0

liuzqt and others added 2 commits September 29, 2026 00:50
… after it is fully deserialized

`AccumulatorV2.readObject` registers the accumulator with the `TaskContext`, which
publishes it where other threads can reach it. That runs at `AccumulatorV2`'s own slot
during deserialization, and Java serialization reads class data superclass-first, so no
subclass field has been assigned yet.

A subclass holding its state in a `val` therefore exposes a null for that window. The
reader that hits it in practice is the executor heartbeater, which calls `isZero` on every
registered accumulator every `spark.executor.heartbeatInterval`. The NPE is thrown on the
`executor-heartbeater` thread and propagates out of a `scheduleAtFixedRate` task, so the
executor stops heartbeating permanently and the driver removes it on heartbeat timeout,
losing anything held only on that executor -- e.g. local checkpoint blocks.

SPARK-20977 worked around this in `CollectionAccumulator` by making its state non-final and
reaching it through a lazily-initializing accessor. Every subclass since has had to
re-derive that workaround, and not all have.

Register in `readResolve` instead, which Java serialization invokes only after every class
in the hierarchy has been read, so the accumulator is never observable half-initialized and
no subclass needs to guard its own state. `readResolve` is `final` for the same reason
`writeReplace` is: unlike the private, per-class `readObject`, it is inherited, so a
subclass defining one would silently stop registration altogether.

Tested with a new case using a subclass that holds its state in a plain `val` with no null
handling: it fails with the NPE on current master and passes with registration moved.
Existing accumulator coverage is unaffected -- core's AccumulatorV2Suite, AccumulatorSuite
and InternalAccumulatorSuite, plus SQLMetricsSuite, AggregatingAccumulatorSuite,
PartitionKeyedAccumulatorSuite and MapperRowCounterSuite in sql.

Co-authored-by: Isaac <no-reply@databricks.com>
…gistration

Deleting `readObject` dropped its `Utils.tryOrIOException` wrapper along with it, which was
not intentional. The wrapper logs the failure and normalizes a non-fatal exception to
`IOException`, and Spark's other `readObject` implementations -- SerializableConfiguration,
SerializableJobConf, Partitioner, SerializableBuffer -- wrap theirs the same way.

Restore it on `readResolve` so this change is only the ordering move, with no incidental
change to how a failure during accumulator deserialization is reported.

Co-authored-by: Isaac <no-reply@databricks.com>

@shrirangmhalgi shrirangmhalgi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Moving registration from readObject to readResolve is the right approach -- it fixes the root cause (publish-before-init) rather than requiring every subclass to independently work around it (SPARK-20977 pattern). Thankyou for the fix.

@cloud-fan cloud-fan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review summary

I reviewed this as a correctness and API-surface change to a public, user-subclassable base class and found no blocking issues. Moving registration from readObject to readResolve correctly defers publication until the whole object is deserialized, establishing a happens-before between the subclass field writes and the point where the heartbeater can read the accumulator via isZero; returning this preserves identity, and dropping defaultReadObject() is safe because Java populates the fields before readResolve runs. Making readResolve final protected mirrors the existing writeReplace convention and is the right call -- it keeps the hook inherited (so subclass instances are still registered) while preventing a subclass from silently overriding it and dropping registration.

I examined two secondary points and concluded neither is a defect. First, readResolve/readObject are Java-serialization plumbing rather than part of the supported implementor contract (which is the abstract methods isZero/copy/copyAndReset/reset/add/merge/value), so sealing readResolve against subclasses -- exactly as writeReplace already is -- breaks no supported external API. Second, adding a non-private readResolve does change the class's default-computed serialVersionUID, but accumulators are only Java-serialized within a single Spark version (task closures and results share bytecode and compute the same UID), and there is no supported cross-version Java-deserialization path for this internal type, so the UID change has no reachable consequence. The new test soundly targets the registration-timing regression (it would NPE under the old ordering) and asserts that the probe actually fired so the check cannot be silently skipped; the readResolve return value is left unasserted here but is exercised pervasively by ordinary task execution.

Findings

0 total: 0 P0, 0 P1, 0 P2, 0 P3.

No findings.

@cloud-fan cloud-fan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Moving TaskContext registration from readObject to readResolve correctly defers publication until the object is fully deserialized, fixing the publish-before-init hazard at the base class rather than per subclass. final protected mirrors the existing writeReplace and is the right call.

@cloud-fan cloud-fan closed this in 5b7d3fe Sep 29, 2026
cloud-fan pushed a commit that referenced this pull request Sep 29, 2026
… after it is fully deserialized

### What changes were proposed in this pull request?

Move the `TaskContext` registration out of `readObject` and into `readResolve`.

`AccumulatorV2` currently registers from `readObject`, which runs at AccumulatorV2's own slot. Java serialization reads hierarchy data **superclass-first**, so no subclass field is assigned yet — while `registerAccumulator` publishes `this` into `TaskMetrics` where other threads can read it. The PR drops `readObject` (beyond registering, all it did was in.defaultReadObject(), which the framework does anyway) and registers from `final protected def readResolve()` instead, which runs only after every class in the hierarchy has been read.

It's `final` for the same reason `writeReplace` already is: `readObject` must be `private` so it's invoked per class, whereas `readResolve` is inherited — a subclass defining one would silently stop registration.

### Why are the changes needed?

Subclasses are published before their state exists, so a subclass holding state in a `val` exposes a `null`. The reader that hits it is the executor heartbeater, calling `isZero` on every registered accumulator. The NPE propagates out of a `scheduleAtFixedRate` task, so `ScheduledThreadPoolExecutor` stops rescheduling — the executor never heartbeats again, the driver removes it on timeout.

Plus the JLS 17.5.3 aspect: a val is a final field written reflectively after publication, so a reader has no guarantee of seeing it and the JIT may fold the read.

SPARK-20977 fixed this locally in CollectionAccumulator; every subclass since has re-derived the workaround and not all have — SetAccumulator in sql/execution/debug still dereferences private val _set from isZero, and EventTimeStatsAccum escapes the NPE only incidentally via Scala's null-safe ==, reporting the wrong answer instead.

### Does this PR introduce _any_ user-facing change?
NO

### How was this patch tested?
New test case

### Was this patch authored or co-authored using generative AI tooling?
Yes. Claude Opus 5.0

Closes #59124 from liuzqt/accumulator-register-after-deserialization.

Authored-by: Ziqi Liu <ziqi.liu@databricks.com>
Signed-off-by: Wenchen Fan <wenchen@databricks.com>
(cherry picked from commit 5b7d3fe)
Signed-off-by: Wenchen Fan <wenchen@databricks.com>
cloud-fan pushed a commit that referenced this pull request Sep 29, 2026
… after it is fully deserialized

### What changes were proposed in this pull request?

Move the `TaskContext` registration out of `readObject` and into `readResolve`.

`AccumulatorV2` currently registers from `readObject`, which runs at AccumulatorV2's own slot. Java serialization reads hierarchy data **superclass-first**, so no subclass field is assigned yet — while `registerAccumulator` publishes `this` into `TaskMetrics` where other threads can read it. The PR drops `readObject` (beyond registering, all it did was in.defaultReadObject(), which the framework does anyway) and registers from `final protected def readResolve()` instead, which runs only after every class in the hierarchy has been read.

It's `final` for the same reason `writeReplace` already is: `readObject` must be `private` so it's invoked per class, whereas `readResolve` is inherited — a subclass defining one would silently stop registration.

### Why are the changes needed?

Subclasses are published before their state exists, so a subclass holding state in a `val` exposes a `null`. The reader that hits it is the executor heartbeater, calling `isZero` on every registered accumulator. The NPE propagates out of a `scheduleAtFixedRate` task, so `ScheduledThreadPoolExecutor` stops rescheduling — the executor never heartbeats again, the driver removes it on timeout.

Plus the JLS 17.5.3 aspect: a val is a final field written reflectively after publication, so a reader has no guarantee of seeing it and the JIT may fold the read.

SPARK-20977 fixed this locally in CollectionAccumulator; every subclass since has re-derived the workaround and not all have — SetAccumulator in sql/execution/debug still dereferences private val _set from isZero, and EventTimeStatsAccum escapes the NPE only incidentally via Scala's null-safe ==, reporting the wrong answer instead.

### Does this PR introduce _any_ user-facing change?
NO

### How was this patch tested?
New test case

### Was this patch authored or co-authored using generative AI tooling?
Yes. Claude Opus 5.0

Closes #59124 from liuzqt/accumulator-register-after-deserialization.

Authored-by: Ziqi Liu <ziqi.liu@databricks.com>
Signed-off-by: Wenchen Fan <wenchen@databricks.com>
(cherry picked from commit 5b7d3fe)
Signed-off-by: Wenchen Fan <wenchen@databricks.com>
@cloud-fan

Copy link
Copy Markdown
Contributor

Merge Summary:

Posted by merge_spark_pr.py

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants