[ 
https://issues.apache.org/jira/browse/SPARK-59418?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59418:
-----------------------------------
    Labels: pull-request-available  (was: )

> Prefer idle executors and isolate repeated OOM retries
> ------------------------------------------------------
>
>                 Key: SPARK-59418
>                 URL: https://issues.apache.org/jira/browse/SPARK-59418
>             Project: Spark
>          Issue Type: Improvement
>          Components: Spark Core
>    Affects Versions: 5.0.0
>            Reporter: Chao Sun
>            Assignee: Chao Sun
>            Priority: Major
>              Labels: pull-request-available
>
> h2. Motivation
> A task that fails with out-of-memory can be retried under the same memory 
> pressure. An executor with free CPU slots may still have several tasks 
> competing for its memory, so moving the retry to another busy executor need 
> not improve its chance of finishing.
> For example, consider an eight-core executor running eight one-CPU tasks. A 
> large partition may fit when fewer tasks compete with it, yet repeatedly fail 
> at full concurrency. If an OOM terminates the executor, the other running 
> tasks also lose their work. Reducing concurrency for the whole application 
> penalizes tasks that do not need additional headroom. Increasing 
> spark.task.cpus reduces concurrency indirectly but does not reserve a 
> proportional amount of memory.
> h2. Proposed approach
> Add an opt-in scheduler recovery policy, disabled by default, that changes 
> placement only after a recognized OOM failure.
> After the first OOM, prefer an eligible idle executor for the retry, with 
> ordinary placement as a fallback. This preference does not reserve the 
> executor: other tasks may subsequently share it.
> After two OOM failures of the same task, temporarily reserve an idle or 
> least-busy compatible executor. Allow existing tasks to finish without 
> preemption, keep new work off that executor, then run the retry alone until 
> its attempt terminates. At most one executor per application is reserved. 
> Bound pending isolation by a configurable timeout (60 seconds by default, 
> measured from the latest OOM); expiry permits ordinary placement and does not 
> interrupt an already-running isolated attempt.
> Keep task CPU requests, executor memory limits, resource profiles, and the 
> task failure budget unchanged. OOM retries may bypass preferred locations but 
> must respect exclusions and CPU/custom-resource requirements. Barrier and 
> pipelined task sets remain outside this recovery policy. Dynamic allocation 
> accounts for the reservation separately so unused reserved slots do not hide 
> ordinary demand, within existing executor limits.
> Use typed OOM signals from task exceptions and executor loss, including 
> JVM/Spark OOM exceptions, app-caused Spark OOM exit codes, and observed 
> Kubernetes executor-container OOMKilled termination. Plain exit code 137 is 
> insufficient. Preserve existing Kubernetes allocator recovery behavior.
> h2. Limits and related work
> This policy reduces competing task activity; it does not estimate per-task 
> memory, increase memory capacity, or guarantee that a partition fits alone. 
> Cached/native memory can remain on idle executors, and unrecognized native 
> failures or OOMs visible only in a restarted container's lastState remain 
> outside the policy. Any production rescue rate or performance benefit 
> requires a controlled enabled/disabled evaluation.
> Related: SPARK-58187 / https://github.com/apache/spark/pull/57329 proposes 
> increasing CPU requests on OOM retries. This proposal instead retains the 
> task resource contract and uses bounded placement/isolation. SPARK-21082 
> discussed memory-aware placement; this proposal does not require memory 
> estimates or rebalance cached data, and explicitly bounds the tradeoff with 
> locality and concurrency.
> Implementation targets current master (5.0.0). Focused validation covers 
> placement, fractional CPU accounting, timeout/cancellation, barrier rollback, 
> dynamic allocation, exception/event-log propagation, and conservative 
> Kubernetes OOM attribution.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to