如何阻止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
相关产品推荐
相关产品推荐

