[SPARK-59845][CORE] Register an accumulator with the TaskContext only after it is fully deserialized - #59124
[SPARK-59845][CORE] Register an accumulator with the TaskContext only after it is fully deserialized#59124liuzqt wants to merge 2 commits into
Conversation
… 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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
… 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>
… 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>
What changes were proposed in this pull request?
Move the
TaskContextregistration out ofreadObjectand intoreadResolve.AccumulatorV2currently registers fromreadObject, which runs at AccumulatorV2's own slot. Java serialization reads hierarchy data superclass-first, so no subclass field is assigned yet — whileregisterAccumulatorpublishesthisintoTaskMetricswhere other threads can read it. The PR dropsreadObject(beyond registering, all it did was in.defaultReadObject(), which the framework does anyway) and registers fromfinal protected def readResolve()instead, which runs only after every class in the hierarchy has been read.It's
finalfor the same reasonwriteReplacealready is:readObjectmust beprivateso it's invoked per class, whereasreadResolveis 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
valexposes anull. The reader that hits it is the executor heartbeater, callingisZeroon every registered accumulator. The NPE propagates out of ascheduleAtFixedRatetask, soScheduledThreadPoolExecutorstops 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