Dask LocalCluster设processes=False调用progress报CancelledError
Dask LocalCluster(processes=False)下调用progress函数抛出asyncio.exceptions.CancelledError问题
- 环境:Dask 2023.9.2
- 问题现象:配置
LocalCluster(processes=False, n_workers=1, threads_per_worker=2)并配合Client执行任务,调用persist后使用progress函数,任务完成后抛出asyncio.exceptions.CancelledError;切换为processes=True配置时程序运行正常。
示例代码
import numpy as np import zarr from dask.distributed import Client, LocalCluster from dask import array as da from dask.distributed import progress def same(x): return x x = np.arange(100000) with LocalCluster(processes=False, n_workers=1, threads_per_worker=2) as cluster, Client(cluster) as client: da_x = da.from_array(x,chunks=(2000,)) da_y = da_x.map_blocks(same) futures = client.persist(da_y) progress(futures,notebook=False) y = da.compute(futures)[0]
报错信息
[########################################] | 100% Completed | 0.3s 2024-04-15 18:37:18,521 - distributed.core - ERROR - Traceback (most recent call last): File "/users/kangl/miniforge3/envs/work/lib/python3.10/site-packages/distributed/utils.py", line 803, in wrapper return await func(*args, **kwargs) File "/users/kangl/miniforge3/envs/work/lib/python3.10/site-packages/distributed/scheduler.py", line 7203, in feed await asyncio.sleep(interval) File "/users/kangl/miniforge3/envs/work/lib/python3.10/asyncio/tasks.py", line 605, in sleep return await future asyncio.exceptions.CancelledError
同类问题参考
在Dask社区发现了一个未解决的同类问题,用户反馈在监控分布式任务进度时遇到与asyncio相关的错误,场景与本次问题高度相似,但目前尚未有官方解决方案。
临时处理方案
- 改用
processes=True的集群配置(已验证可规避该错误) - 在
progress调用后添加短暂延迟,让异步监控任务有足够时间正常退出,例如:import time progress(futures,notebook=False) time.sleep(0.5) - 捕获该错误(任务实际已完成,此错误不影响最终结果),例如:
import asyncio try: progress(futures,notebook=False) except asyncio.exceptions.CancelledError: pass
内容的提问来源于stack exchange,提问作者Kang Liang
相关产品推荐
相关产品推荐

