[
https://issues.apache.org/jira/browse/SPARK-59561?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-59561:
-----------------------------------
Labels: pull-request-available (was: )
> ExecutorAllocationManager permanently stalls if a SparkListenerStageSubmitted
> event is dropped from the executorManagement queue
> --------------------------------------------------------------------------------------------------------------------------------
>
> Key: SPARK-59561
> URL: https://issues.apache.org/jira/browse/SPARK-59561
> Project: Spark
> Issue Type: Bug
> Components: Spark Core
> Affects Versions: 4.0.1, 5.0.0
> Reporter: dongjunhwang
> Priority: Major
> Labels: pull-request-available
>
> h3. Problem
> SPARK-58935 documented that {{LiveListenerBus}} delivers the same logical
> event to several independent, separately-capacity-limited queues, and that
> {{ExecutorAllocationManager}}'s listener (registered on the
> {{executorManagement}} queue via {{addToManagementQueue}},
> [LiveListenerBus.scala#L70|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/scheduler/LiveListenerBus.scala#L70])
> can permanently miss a {{SparkListenerStageSubmitted}} event if that queue
> is full when the event is posted ({{AsyncEventQueue.post}},
> [AsyncEventQueue.scala#L168|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/scheduler/AsyncEventQueue.scala#L168]).
> {{ExecutorAllocationManager}}'s "executors needed" calculation
> ({{maxNumExecutorsNeededPerResourceProfile}},
> [ExecutorAllocationManager.scala#L416|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/ExecutorAllocationManager.scala#L416])
> is driven entirely by bookkeeping populated in {{onStageSubmitted}}
> ([ExecutorAllocationManager.scala#L820|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/ExecutorAllocationManager.scala#L820]).
> If that event is dropped, this manager never learns the stage exists and
> permanently omits it from the calculation for the stage's entire lifetime --
> there is no periodic resync. Meanwhile the Spark UI/REST API, backed by
> {{AppStatusStore}} via the separate {{appStatus}} queue
> ({{addToStatusQueue}},
> [LiveListenerBus.scala#L75|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/scheduler/LiveListenerBus.scala#L75]),
> can continue to correctly show the stage as {{RUNNING}}, since a drop on one
> queue does not imply a drop on the other.
> SPARK-58935 fixed this only partially, by documenting the existing
> {{queue.executorManagement.numDroppedEvents}} metric so operators can at
> least detect and alert on the condition
> ([PR|https://github.com/apache/spark/pull/58216]). The underlying
> permanent-stall behavior itself -- called out explicitly as unaddressed in
> that PR's description -- is still open. This ticket is that follow-up.
> h3. Reproduction
> Deterministically reproduced by constructing a {{LiveListenerBus}} with
> {{spark.scheduler.listenerbus.eventqueue.executorManagement.capacity=1}},
> occupying the queue's single dispatch thread with a blocking listener, and
> posting a {{SparkListenerStageSubmitted}} event for a multi-task stage once
> the queue is full. The event is dropped, and
> {{ExecutorAllocationManager.maxNumExecutorsNeededPerResourceProfile}} returns
> {{0}} for that stage even though its tasks are pending, while
> {{AppStatusStore.activeStages()}} (fed by the separate, unconstrained
> {{appStatus}} queue) correctly shows the stage as active.
> h3. Proposed fix
> A lightweight, best-effort reconciliation on top of
> {{ExecutorAllocationManager}}'s existing ~100ms polling loop ({{schedule()}},
> [ExecutorAllocationManager.scala#L460|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/ExecutorAllocationManager.scala#L460]):
> * {{schedule()}} cheaply checks {{LiveListenerBus}}'s existing
> {{executorManagement}} drop counter (a single {{Long}} comparison) and, only
> when it has increased, reconciles.
> * Reconciliation diffs {{AppStatusStore.activeStages()}}
> ([AppStatusStore.scala#L171|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/status/AppStatusStore.scala#L171])
> against the manager's own bookkeeping and registers any stage
> {{AppStatusStore}} knows about but the manager doesn't, through the same code
> path as a normal {{onStageSubmitted}}. A recovered stage's ongoing
> pending/running counts are then read fresh from {{AppStatusStore}} (via
> {{AppStatusStore.stageAttempt}},
> [AppStatusStore.scala#L226|https://github.com/apache/spark/blob/9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150/core/src/main/scala/org/apache/spark/status/AppStatusStore.scala#L226])
> rather than reconstructing the manager's normal per-task-index bookkeeping,
> for as long as it remains "recovered," until it completes.
> * The same congestion that drops a stage's {{StageSubmitted}} could later
> also drop its {{StageCompleted}}; reconciliation also sweeps recovered
> attempts whose ground-truth status has since become terminal and cleans them
> up the same way a real {{StageCompleted}} would, so this can't leak for the
> remaining lifetime of the application.
> This intentionally does not introduce a new pub/sub framework analogous to
> K8s's {{ExecutorPodsSnapshotsStore}} (there is only one consumer here), and
> does not attempt a locality-optimized placement for recovered stages (their
> locality preferences aren't available from {{AppStatusStore}}'s lightweight
> {{StageData}}; the goal is a correct executor count, not optimal placement).
> Known limitation: a task that was speculatively submitted but not yet started
> during the drop window can't be recovered, since {{AppStatusStore}}'s
> lightweight path has no ground truth for "requested but not yet running"
> speculative tasks. This self-heals once the task actually starts and is
> picked up by normal running-task accounting.
> I have a working patch and tests and will open a PR shortly.
> h3. Related work
> * SPARK-32597 (still open, unassigned) previously identified this general
> class of problem and proposed a more invasive {{VariableLinkedBlockingQueue}}
> approach ([PR #29413|https://github.com/apache/spark/pull/29413], closed
> unmerged in 2020). {{mridulm}} noted a simpler fix (manually setting queue
> capacity) already existed; {{cloud-fan}} noted no queue-capacity approach can
> fully eliminate drops; {{tgravescs}} separately suggested the more complete
> fix is integrating dynamic allocation directly into the core scheduler rather
> than depending on the listener bus at all -- that hasn't been implemented in
> the 5 years since. This proposal does not attempt that larger redesign; it
> takes a complementary angle, detecting drops cheaply after the fact and
> repairing their effect on dynamic allocation specifically, using ground truth
> Spark already maintains elsewhere.
> * SPARK-58935 (Resolved, Fixed) documented the specific manifestation and
> added the diagnostic metric mentioned above; this ticket implements the "full
> fix" its PR description called out as future work.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]