如何正确嵌套dask.delayed函数以获得预期执行结果?
我正在学习Dask,创建了一个延迟管道示例,依赖关系为:baz依赖bar,bar依赖foo。我希望这三个函数都作为delayed任务执行,且在仪表板上显示为独立任务。
我的代码如下:
import dask import distributed @dask.delayed def foo(data:str) -> str: return f'foo[{data}]' @dask.delayed def bar(data:str) -> str: return f'bar[{data}]' @dask.delayed def baz(data:str) -> str: f = foo(data) b = bar(f) return f'baz[{b}]' baz_task = baz('hello world') client = distributed.Client() future = client.compute(baz_task) result = future.result() print(result)
运行后返回结果为:
baz[Delayed('bar-4c8a0ec0-1d87-43e7-a8af-fee33a9fae3d')]
而我预期的结果是:
baz[bar[foo[hello world]]]
我尝试过两种方法:
- 移除foo和bar的@dask.delayed装饰器:导致任务图中只有一个任务,无法在仪表板跟踪流程;
- 在baz函数内部调用Delayed对象:返回延迟的apply对象,仍未解决问题。
我的代码似乎符合官方文档的装饰器示例,请问我哪里出错了?
问题原因与解决方案
问题出在被@dask.delayed装饰的函数内部,直接使用延迟对象会保留其Delayed实例的字符串表示,而不会自动触发计算并获取结果。因为baz本身是延迟任务,它内部调用foo和bar得到的是Delayed对象,而非实际计算后的字符串,所以格式化字符串时会把Delayed对象的默认打印值拼进去。
要解决这个问题,需要让baz接收bar的计算结果作为参数,将依赖关系移到函数外部构建,让每个延迟任务的输入是前一个任务的输出:
import dask import distributed @dask.delayed def foo(data:str) -> str: return f'foo[{data}]' @dask.delayed def bar(data:str) -> str: return f'bar[{data}]' @dask.delayed def baz(data:str) -> str: return f'baz[{data}]' # 在外部构建依赖链 foo_task = foo('hello world') bar_task = bar(foo_task) baz_task = baz(bar_task) client = distributed.Client() future = client.compute(baz_task) result = future.result() print(result)
这样修改后,三个函数都会作为独立的延迟任务出现在Dask仪表板中,最终输出结果就是预期的baz[bar[foo[hello world]]]。
原代码失效的核心原因
当baz被@dask.delayed装饰后,它的内部逻辑会被延迟执行。在baz内部调用foo(data)和bar(f)时,得到的是Delayed对象,但这些对象并不会在baz的延迟计算过程中自动解析为实际结果——Dask只会把baz的函数体作为一个单独任务,其中的Delayed对象会被当作普通参数处理,所以最终返回的是包含Delayed实例字符串的结果。
而把依赖链移到外部后,Dask会正确构建任务依赖图:foo完成后执行bar,bar完成后执行baz,每个任务都能拿到前一个任务的实际计算结果,同时三个任务在仪表板上都是独立可跟踪的。
内容的提问来源于stack exchange,提问作者Steve Lorimer

