Jonathan Lamberts created SPARK-59451:
-----------------------------------------
Summary: When Heartbeater calls isZero on partially-deserialized
Accumulator, it can cause an NPE which kills the heartbeat thread
Key: SPARK-59451
URL: https://issues.apache.org/jira/browse/SPARK-59451
Project: Spark
Issue Type: Bug
Components: Spark Core
Affects Versions: 4.2.0
Reporter: Jonathan Lamberts
There is a race condition that can occur when the heartbeat thread calls isZero
on an accumulator while it's mid-deserialization. I originally encountered this
via PartitionKeyedAccumulator, but theoretically it could happen in any
accumulator whose isZero method depends on a subclass-defined field (and
actually has happened before in CollectionAccumulator).
Basically, since Java deserializes class hierarchies top-down, there is a brief
window when an accumulator is registered but its fields are not fully
populated. In the case of PartitionKeyedAccumulator, the relevant field is a
hashmap named byPartition. If the heartbeat fires during this window, it calls
isZero, which causes a NPE since the hashmap is still null. The heartbeat
wrapper logs this then reraises, which stops future heartbeats. This means the
executor keeps working as a zombie until the heartbeat timeout is reached, then
the driver kills it.
The stack trace (this is from DBR so the line numbers won't match up):
Uncaught exception in thread executor-heartbeater, stack trace:
java.lang.NullPointerException: Cannot invoke "ConcurrentHashMap.isEmpty()"
because the return value of
"org.apache.spark.sql.util.PartitionKeyedAccumulator.byPartition()" is null
at
org.apache.spark.sql.util.PartitionKeyedAccumulator.isZero(PartitionKeyedAccumulator.scala:51)
at
org.apache.spark.executor.Executor.getHeartbeatAccumUpdatesForTask(Executor.scala:2586)
at org.apache.spark.executor.Executor.reportHeartBeat(Executor.scala:2648)
at org.apache.spark.Heartbeater$$anon$1.run(Heartbeater.scala:46)
{{}}
{{}}
This is somewhat hard to reproduce unless you run a lot of tasks, but in our
prod job with ~4 million tasks that contain the serialized accumulator it's
happening 6-8 times per run. The mechanism itself is unit testable by forcing
an isZero call mid-deserialization.
Rather than handling the NPE inside of the accumulator, I think it's probably
better to defer registering the accumulator entirely until after
deserialization ends. If there's some reason we don't want to do this, I can
put in a fix for PartitionKeyedAccumulator specifically.
There's also a separate question of whether the heartbeater itself should die
on a non-fatal exception, but I don't know if there are other implications
there so leaving it out of scope for now.
I'm happy to open a PR to fix this issue.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]