Prefect 2 + Dask:提交任务未占用GPU资源问题求助
问题:Prefect 2 + Dask GPU资源调度失效导致Worker过载崩溃
目标
让Prefect 2创建的Dask任务占用GPU资源,避免Worker因任务过载崩溃。
已确认配置
- 每个Dask Worker已配置
GPU=1的资源配额 - 通过Dask仪表盘验证所有Worker的GPU资源显示正常
问题现象
通过Prefect 2运行任务时,GPU资源未被标记为已占用,多个任务同时分配到同一Worker,最终导致Worker过载崩溃。
当前代码实现
import requests from prefect import flow, task, get_run_logger from prefect_dask.task_runners import DaskTaskRunner import dask @task def UpscaleFrames(FramesToUpscale): # 执行CUDA相关的帧放大操作 return @flow(task_runner=DaskTaskRunner(address="tcp://tower:8786")) def Upscale(): for file in GetVideoFiles("/videos"): while (frames_found): FramesToUpscale = GetFramesToUpscale() with dask.annotate(resources={'GPU': 1}): UpscaleFrames.submit(FramesToUpscale)
环境版本信息
Version: 2.3.2 API version: 0.8.0 Python version: 3.10.6 Git commit: 6e931ee9 Built: Tue, Sep 6, 2022 12:36 PM OS/Arch: linux/x86_64 Profile: default Server type: hosted
解决方案
问题出在dask.annotate上下文无法传递到Prefect任务的底层Dask调度逻辑中,Prefect的task.submit方法需要通过task_runner_kwargs直接指定资源要求。
修改后的代码
import requests from prefect import flow, task, get_run_logger from prefect_dask.task_runners import DaskTaskRunner import dask @task def UpscaleFrames(FramesToUpscale): # 执行CUDA相关的帧放大操作 return @flow(task_runner=DaskTaskRunner(address="tcp://tower:8786")) def Upscale(): for file in GetVideoFiles("/videos"): while (frames_found): FramesToUpscale = GetFramesToUpscale() # 直接通过task_runner_kwargs传递GPU资源要求 UpscaleFrames.submit(FramesToUpscale, task_runner_kwargs={"resources": {"GPU": 1}})
额外验证步骤
- 任务提交后,查看Dask仪表盘,确认每个
UpscaleFrames任务运行时会占用对应Worker的GPU资源 - 确认Dask Worker启动时确实指定了资源参数:
dask-worker tcp://tower:8786 --resources GPU=1
- 考虑升级Prefect及prefect-dask到较新版本(当前使用的2.3.2是2022年的旧版本,新版本对Dask资源调度的兼容性更好)
内容的提问来源于stack exchange,提问作者Zap
相关产品推荐
相关产品推荐

