如何实现不占用Dask Worker的休眠等待函数?
解决Dask等待外部进程时不阻塞Worker的问题
嘿,这个场景我太熟悉了!原来那种在Worker上循环time.sleep的写法,会把Worker线程死死占住,这段时间里它啥别的任务都干不了,简直是资源浪费。下面给你几个实用的方案,都能做到等待的时候不占用Worker资源:
核心思路
本质上就是别让Worker抱着任务一直休眠——把“等待+检查”拆成一个个短任务,每次检查完如果文件还没 ready,就让任务结束,让Dask调度器过一会儿再重新调度这个检查任务。这样中间的等待时间里,Worker就能去处理别的活儿了。
方案一:用Tenacity库做重试(推荐,代码更简洁)
Tenacity是个专门处理重试逻辑的库,用它来控制检查频率和超时时间特别方便,先装一下:pip install tenacity
import time from pathlib import Path from dask import delayed from tenacity import retry, stop_after_delay, wait_fixed, retry_if_result # 定义重试触发条件:文件不存在就重试 def file_not_ready(filename): return not Path(filename).exists() @retry( stop=stop_after_delay(600), # 最多等10分钟 wait=wait_fixed(0.1), # 每隔0.1秒检查一次 retry=retry_if_result(file_not_ready) ) def wait_for_file(filename='some_filename'): if Path(filename).exists(): return filename # 返回False触发重试 return False # 用delayed包装一下就能放进Dask任务流里了 file_exists = delayed(wait_for_file)('some_filename')
方案二:手动实现重试逻辑(不用第三方库)
如果不想加依赖,自己写个简单的重试逻辑也可以,通过抛出特定异常让Dask重新调度任务:
import time from pathlib import Path from dask import delayed # 自定义一个异常,标记文件还没准备好 class FileNotReady(Exception): pass def wait_for_file(filename='some_filename', max_wait_time=600, check_interval=0.1): start_time = time.time() if Path(filename).exists(): return filename # 超时就抛出TimeoutError if time.time() - start_time > max_wait_time: raise TimeoutError(f"等了太久啦,文件{filename}还是没出现!") # 没超时也没找到文件,抛出异常让Dask稍后重试 raise FileNotReady(f"文件{filename}还没好,稍后再检查~") # 配置Dask允许重试这个自定义异常,设置重试次数和间隔 from dask.config import set_options with set_options(retries=1000, retry_delay=0.1): file_exists = delayed(wait_for_file)('some_filename')
方案三:用异步任务(适合分布式场景)
如果用的是Dask分布式集群,异步写法更高效——用asyncio.sleep代替time.sleep,异步休眠不会占用Worker线程,Worker可以同时处理其他异步任务:
import asyncio import time from pathlib import Path from dask.distributed import Client, as_completed async def wait_for_file_async(filename='some_filename', max_wait_time=600, check_interval=0.1): start_time = time.time() while True: if Path(filename).exists(): return filename if time.time() - start_time > max_wait_time: raise TimeoutError(f"超时!文件{filename}未在{max_wait_time}秒内出现") # 异步休眠,不阻塞Worker线程 await asyncio.sleep(check_interval) # 提交异步任务到集群 client = Client() future = client.submit(wait_for_file_async, 'some_filename') # 等待任务完成 for completed in as_completed([future]): result = completed.result() print(f"搞定!文件已就绪:{result}")
这几个方案里,我个人推荐用Tenacity的那个,代码最简洁,重试逻辑也很清晰,不容易出错。异步的方案适合大规模分布式集群,资源利用率最高。
内容的提问来源于stack exchange,提问作者mvn
相关产品推荐
相关产品推荐

