You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

dask distributed中delayed函数返回列表的正确使用及问题解决

Dask Distributed 问题解决方案

问题1:len(tasks)输出不符合预期

原因

do_tasks是普通Python函数,没有被@dask.delayed装饰,调用时会立即执行,返回的是(list_x, list_y)组成的二元元组,元组本身长度天然就是2。你要获取对应日期的任务数量,需要取元组内部list_x或list_y的长度,两个列表都保存了5个对应日期的延迟任务对象。

修正方法

将代码中的

print(len(tasks))

修改为

print(len(list_x))

即可输出预期值5。


问题2:避免重复执行do_tasks

原因

Dask默认不会自动缓存中间计算结果,两次单独调用compute时,都会从头遍历整个依赖链执行任务,导致公共的do_tasks逻辑被重复计算两次。

修正方案

有两种常用实现方式:

方案1:预持久化公共中间结果

调用persist方法提前把tasks的计算结果存储到集群内存中,后续两次计算直接读取缓存结果,无需重复执行公共逻辑:

import dask
import operator

@dask.delayed
def do_something(date):
    # 此处为你的原有逻辑,x、y为实际返回值
    x = f"x_{date}"
    y = f"y_{date}"
    return x, y

get_item0 = dask.delayed(operator.itemgetter(0))
get_item1 = dask.delayed(operator.itemgetter(1))

def handle_x(list_x):
    print(len(list_x))
    return list_x

def handle_y(list_y):
    print(len(list_y))
    return list_y

def do_tasks():
    list_x, list_y = [], []
    dates = [20210101, 20210102, 20210103, 20210104, 20210105]
    for date in dates:
        result = do_something(date)
        x = get_item0(result)
        y = get_item1(result)
        list_x.append(x)
        list_y.append(y)
    return list_x, list_y

if __name__ == "__main__":
    # 注意原代码的拼写错误:Distributed首字母应为小写
    with dask.distributed.Client() as dask_client:
        tasks = do_tasks()
        list_x = get_item0(tasks)
        list_y = get_item1(tasks)

        print(len(list_x))

        # 持久化公共中间结果到集群内存
        tasks_persisted = dask_client.persist(tasks)
        list_x_persisted = get_item0(tasks_persisted)
        list_y_persisted = get_item1(tasks_persisted)

        # 两次计算直接复用缓存结果
        dask_client.compute(dask.delayed(handle_x)(list_x_persisted)).result()
        dask_client.compute(dask.delayed(handle_y)(list_y_persisted)).result()

方案2:合并计算任务单次提交

把两个处理函数的延迟任务合并为一个列表单次提交,Dask会自动识别依赖链中的公共部分,只执行一次公共逻辑:

if __name__ == "__main__":
    with dask.distributed.Client() as dask_client:
        tasks = do_tasks()
        list_x = get_item0(tasks)
        list_y = get_item1(tasks)

        print(len(list_x))

        # 合并两个任务单次提交
        fx = dask.delayed(handle_x)(list_x)
        fy = dask.delayed(handle_y)(list_y)
        res_x, res_y = dask_client.compute([fx, fy]).result()

补充说明:你原代码中的dask.lazy是旧版本接口,现在统一使用dask.delayed即可,另外dask.Distributed.Client存在拼写错误,正确写法为dask.distributed.Client。


内容的提问来源于stack exchange,提问作者abisko

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 07:57:02