You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何阻止Dask Client提交的函数在Worker因CancelledError(如OOM)被杀后重启

解决Dask Worker因OOM被杀时任务自动重试的问题

问题核心

你遇到的情况本质是两类重试逻辑的控制范围不同:

  • Client.submit(retries=N)仅负责任务在Worker上执行时主动抛出业务异常的重试次数,默认值0确实会让这类异常直接终止任务
  • 而Worker因OOM被杀导致的任务丢失,属于集群层面的不可预期失败,这类重试由Dask调度器的全局策略控制,默认会重试4次,不受submit的retries参数影响

解决方案

方法1:全局禁用调度器重试

创建Client时通过defaults参数直接修改调度器的重试配置,对所有任务生效:

from dask.distributed import Client
import os

def oom_script():
    def generate_data() -> bytes:
        return os.urandom(10) + b":-) " * (100_000_000 // 4)
    oom_list = []
    while True:
        oom_list.append(generate_data())

# 初始化Client时设置调度器重试次数为0
client = Client(defaults={"distributed.scheduler.retries": 0})
client.submit(oom_script)

方法2:仅针对单个任务禁用重试

如果不想全局修改配置,可通过临时上下文覆盖调度器配置,仅对指定任务生效:

from dask.distributed import Client
import dask.config
import os

def oom_script():
    def generate_data() -> bytes:
        return os.urandom(10) + b":-) " * (100_000_000 // 4)
    oom_list = []
    while True:
        oom_list.append(generate_data())

client = Client()

# 临时设置调度器重试为0,仅在此上下文内的任务生效
with dask.config.set({"distributed.scheduler.retries": 0}):
    future = client.submit(oom_script)

原理补充

Dask调度器默认会对「任务丢失(Worker被杀、网络中断等)」这类非业务异常进行重试,默认重试次数由distributed.scheduler.retries配置项控制(默认值4)。而Client.submit的retries参数仅作用于任务执行中主动抛出的异常(比如你提供的exception_script示例),因此无法阻止Worker被杀后的重试行为。

内容的提问来源于stack exchange,提问作者Petr Synek

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 06:50:08