把 Python 工作流转换为 Dask Delayed 任务图

把 Python 工作流转换为 Dask Delayed 任务图

一段 Python 程序可以把文件读取、计算与汇总写成连续的函数调用;当输入变大或工作可并行时,Dask Delayed 允许将函数调用记录为任务图,再由调度器决定执行顺序与并行度。关键变化是先描述依赖,再执行,而非每一步都立即计算。

延迟调用并在依赖完整后计算

用 dask.delayed 包装普通函数,调用后获得代表未来结果的延迟对象。把该对象传给后续延迟函数,依赖关系就进入图中;最后调用 dask.compute,一次执行完整图并返回具体结果。

三个延迟读取与计算任务汇入一个汇总任务,compute 在图构建完成后统一执行
图:只有在最终汇总的依赖关系明确后,调度器才能安排整张任务图。
from dask import delayed

@delayed
def read_value(path):
    return read_file(path)

@delayed
def analyze(value):
    return expensive_analysis(value)

parts = [analyze(read_value(path)) for path in paths]
result = delayed(sum)(parts).compute()

这段伪例假设已有安全、正确的 read_file 与 expensive_analysis 函数。不要把整张延迟对象转成字符串或在构图时意外解包;如需逐个计算或共享中间结果,优先考虑统一提交相关输出,避免重复执行同一子图。

明确依赖,控制任务粒度

任务应通过参数和返回值显式传递依赖。若任务暗中读取可变全局变量、当前目录或远程服务状态,调度器无法从图看出这些关系,分布式执行还可能在不同 worker 得到不同结果。对副作用操作(写文件、发请求、更新数据库)要特别谨慎,因为重试或重复调度可能再次产生副作用;应设计幂等操作、唯一输出路径和明确的提交边界。

过多极小任务会让调度与序列化开销超过计算本身。把一段成本足够高且独立的工作作为任务,避免每个简单算术都变成节点。反过来,过大的任务会减少可并行度与故障恢复能力,需根据数据大小、计算时长和可用内存调整分块。

与 Dask 集合协作

Dask Array、DataFrame 等集合已有自己的分块任务图。不要先把大型集合转换为一个巨大的普通 Python 对象再交给 Delayed,这会破坏分块优势并可能耗尽内存。可把分区映射为 Delayed 单元,或在分块操作中处理数据;只在最终结果确实适合内存时才计算并物化。

调度器也影响共享数据语义。模块级全局对象在本地线程调度器中可能被线程共同访问,但进程或分布式调度器通常需要序列化或把数据发送给其他 worker;不能把“全局变量共享”当作跨调度器通用保证。尽量用图节点的输入输出表达数据流,让任务在不同执行后端下保持可理解。

来源:Dask 官方《Delayed》与《Best Practices》文档。Dask API 和调度器行为可能随版本更新;请对照项目实际版本及数据规模选择调度器和任务粒度。

Dask 来源许可

本文涉及的 Dask 项目文档与示例来源按上游 Dask 仓库的 BSD 3-Clause 声明标注。保留版权和原许可信息;本文中文译写另行授权。上游声明的版权人为 Anaconda, Inc. 与贡献者,完整 BSD 3-Clause 文本如下。

BSD 3-Clause License

Copyright (c) 2014, Anaconda, Inc. and contributors
All rights reserved.

Redistribution and use in source and binary forms, with or without
modification, are permitted provided that the following conditions are met:

* Redistributions of source code must retain the above copyright notice, this
  list of conditions and the following disclaimer.
* Redistributions in binary form must reproduce the above copyright notice,
  this list of conditions and the following disclaimer in the documentation
  and/or other materials provided with the distribution.

* Neither the name of the copyright holder nor the names of its
  contributors may be used to endorse or promote products derived from
  this software without specific prior written permission.
THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE
FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL
DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR
SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER
CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY,
OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容