[
https://issues.apache.org/jira/browse/SPARK-58908?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-58908:
-----------------------------------
Labels: pull-request-available (was: )
> On Kubernetes, heartbeat-timeout executor loss is reported as ExecutorKilled
> so task failures are not counted
> -------------------------------------------------------------------------------------------------------------
>
> Key: SPARK-58908
> URL: https://issues.apache.org/jira/browse/SPARK-58908
> Project: Spark
> Issue Type: Bug
> Components: Kubernetes
> Affects Versions: 4.2.0
> Reporter: Abhinav Battu
> Priority: Major
> Labels: pull-request-available
>
> h3. Summary
> On Kubernetes, when an executor is removed because its heartbeat timed out,
> the loss is
> recorded as {{ExecutorKilled}} instead of {{ExecutorProcessLost}}.
> {{ExecutorKilled}} is
> classified as not caused by the application, so tasks running on that
> executor do not count
> towards {{spark.task.maxFailures}} and the job never fast-fails. On YARN the
> same scenario
> records {{ExecutorProcessLost}} and the job fails after the retry limit.
> h3. Mechanism
> {{HeartbeatReceiver.expireDeadHosts}} handles a timeout in two steps:
> # {{sc.killAndReplaceExecutor(id)}}, which calls
> {{killExecutors(..., countFailures = true, force = true)}}.
> {{CoarseGrainedSchedulerBackend.killExecutors}} records
> {{executorsPendingToRemove(id) = !countFailures}}, i.e. {{false}}, meaning
> "the driver killed
> this, but the failures must still count".
> # {{driverEndpoint.send(RemoveExecutor(id, ExecutorProcessLost("Executor
> heartbeat timed out
> after N ms")))}}.
> {{KubernetesClusterSchedulerBackend.doKillExecutors}} runs inside step 1 and
> eagerly reports
> its own loss reason for every executor:
> {code:scala}
> executorIds.foreach { id =>
> removeExecutor(id, ExecutorKilled)
> }
> {code}
> Because {{adjustTargetNumExecutors = false}}, {{killExecutors}} resolves
> {{Future.successful(true).flatMap(...)}} on {{ThreadUtils.sameThread}}, so
> {{doKillExecutors}} runs synchronously on the calling thread. The
> {{ExecutorKilled}} message
> is therefore enqueued to the driver endpoint before step 2 sends its message.
> {{DriverEndpoint.removeExecutor}} then computes:
> {code:scala}
> val killedByDriver =
> executorsPendingToRemove.remove(executorId).getOrElse(false)
> if (killedByDriver) ExecutorKilled else ... else reason
> {code}
> {{killedByDriver}} is {{false}} (step 1 stored {{false}}), so the reason
> supplied by the
> caller is used - and that caller is the Kubernetes backend, which supplied
> {{ExecutorKilled}}. The executor is removed from {{executorDataMap}}, so the
> later
> {{ExecutorProcessLost}} message from step 2 falls into the {{case None}}
> branch and has no
> effect.
> {{TaskSetManager}} then maps the reason:
> {code:scala}
> case ExecutorKilled | ExecutorDecommission(_, _) => false // exitCausedByApp
> {code}
> so {{ExecutorLostFailure.countTowardsTaskFailures}} is {{false}} and the
> failures are not
> counted.
> h3. This is deterministic, not a race
> The eager removal does not merely sometimes win. {{doKillExecutors}}
> completes synchronously
> before step 2 runs, so {{ExecutorKilled}} always wins.
> h3. Impact
> On Kubernetes, repeated heartbeat timeouts never advance the task failure
> count, so a job
> retries indefinitely instead of failing after {{spark.task.maxFailures}}
> attempts. YARN is
> unaffected because its {{doKillExecutors}} does not report a loss reason;
> Kubernetes is the
> only backend that does.
> h3. Reproduction
> Reproduced on current master in {{HeartbeatReceiverSuite}} by expiring an
> executor against
> two backends and capturing the reason delivered to
> {{TaskSchedulerImpl.executorLost}}:
> {code}
> control (no eager removal): Executor heartbeat timed out after 240000 ms ->
> counts
> k8s-like (eager removal): Executor killed by driver. ->
> does not count
> {code}
> The existing {{KubernetesClusterSchedulerBackendSuite}} "Kill executors" test
> independently
> asserts that the real backend sends {{RemoveExecutor(id, ExecutorKilled)}}.
> h3. Note
> {{countFailures = true}} has exactly one caller in the codebase
> ({{SparkContext.killAndReplaceExecutor}}, called only from
> {{HeartbeatReceiver.expireDeadHosts}}), so the affected path is precisely the
> heartbeat
> timeout.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]