如何在Kedro中实现节点A、B、C并行执行后再运行节点D并保持代码整洁?
解决Kedro节点并行与顺序执行的问题
当然可以!Kedro的Pipeline调度机制天生就支持这种场景,既能让资源消耗低但耗时久的A、B、C并行执行以节省时间,又能确保占用高内存的D等它们全部完成再启动,同时还能保持代码结构的整洁性。
核心思路
Kedro是基于节点的依赖关系来调度执行顺序的:
- 如果多个节点之间没有互相依赖的输入/输出,Kedro会自动并行执行它们;
- 如果一个节点依赖其他节点的输出(或者通过显式声明依赖),Kedro会等待所有依赖节点完成后再启动该节点。
根据你的需求,我们可以通过两种方式实现,取决于节点D是否需要A、B、C的输出数据:
情况1:D需要A、B、C的输出数据
如果D的逻辑本身需要用到bar_a、bar_b、bar_c这些结果,直接把它们加入D的输入列表即可。Kedro会自动识别依赖关系,先并行跑A、B、C,等它们都生成输出后再启动D:
from kedro.pipeline import Pipeline, node from .nodes import * def foo(): # 单独定义A、B、C节点,逻辑更清晰 node_a = node(a, inputs=["train_x", "test_x"], outputs=dict(bar_a="bar_a"), name="A") node_b = node(b, inputs=["train_x", "test_x"], outputs=dict(bar_b="bar_b"), name="B") node_c = node(c, inputs=["train_x", "test_x"], outputs=dict(bar_c="bar_c"), name="C") # D节点输入包含A、B、C的输出,自动建立依赖 node_d = node( d, inputs=["train_x", "test_x", "bar_a", "bar_b", "bar_c"], outputs=dict(bar_d="bar_d"), name="D" ) return Pipeline([node_a, node_b, node_c, node_d])
情况2:D不需要A、B、C的输出数据
如果D完全不需要bar_a、bar_b、bar_c,只是单纯需要等待A、B、C执行完成,就可以用节点的dependencies参数显式声明依赖关系,这样代码更直观,也不会传入无用的输入:
from kedro.pipeline import Pipeline, node from .nodes import * def foo(): node_a = node(a, inputs=["train_x", "test_x"], outputs=dict(bar_a="bar_a"), name="A") node_b = node(b, inputs=["train_x", "test_x"], outputs=dict(bar_b="bar_b"), name="B") node_c = node(c, inputs=["train_x", "test_x"], outputs=dict(bar_c="bar_c"), name="C") # 显式指定D依赖A、B、C三个节点的完成,无需传入它们的输出 node_d = node( d, inputs=["train_x", "test_x"], outputs=dict(bar_d="bar_d"), name="D", dependencies=["A", "B", "C"] # 通过节点名称声明依赖 ) return Pipeline([node_a, node_b, node_c, node_d])
效果说明
两种写法都能达成你的需求:
- A、B、C之间没有互相依赖,Kedro会自动并行执行它们,充分利用计算资源;
- D必须等待A、B、C全部执行完成后才会启动,彻底避免了内存资源冲突的问题;
- 代码结构清晰,节点单独定义后再组合成Pipeline,可读性和维护性都很好。
内容的提问来源于stack exchange,提问作者João Areias
相关产品推荐
相关产品推荐

