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

如何让Kedro pipeline的输入DataFrame可根据运行模式配置?

实现Kedro Pipeline输入动态适配运行模式的方案

核心实现思路

通过Kedro原生的多环境配置+数据集别名机制实现无侵入的动态切换,不需要修改Pipeline核心逻辑:

  • Pipeline内部只绑定通用的数据集别名,不关联具体的数据源路径
  • 独立运行模式下,将别名映射到本地测试用的CSV文件
  • 集成运行模式下,将同一个别名映射到前置Pipeline的输出数据集(可选择磁盘持久化文件或内存数据集)

具体配置步骤

1. 定义基础目录配置

在conf/base/catalog.yml中定义全局通用的数据集,PL1的输入别名不做默认绑定,留到不同环境配置中重写:

# conf/base/catalog.yml
# PL0的输出数据集,独立运行PL0时输出到该路径,集成模式下可直接作为PL1输入
pl0_output_dataset:
  type: pandas.CSVDataSet
  filepath: data/02_intermediate/pl0_output.csv

# PL1的输出数据集
pl1_output_dataset:
  type: pandas.CSVDataSet
  filepath: data/03_primary/pl1_output.csv

2. 配置独立运行模式的数据源

在conf/local/catalog.yml(本地开发默认加载的环境配置)中,将PL1的输入别名映射到测试CSV文件:

# conf/local/catalog.yml 仅本地独立运行时生效
pl1_input_dataset:
  type: pandas.CSVDataSet
  filepath: data/01_raw/pl1_test_input.csv # 独立运行PL1时的测试输入文件

3. 配置集成运行模式的数据源

在conf/prod/catalog.yml(生产环境配置)中,将PL1的输入别名直接映射到PL0的输出:

# conf/prod/catalog.yml 生产串联运行时生效
pl1_input_dataset:
  type: pandas.CSVDataSet
  filepath: data/02_intermediate/pl0_output.csv # 直接复用PL0的输出文件

# 不需要持久化中间数据时可以替换为内存数据集,性能更高:
# pl1_input_dataset:
#   type: pandas.MemoryDataSet

Pipeline代码示例

Pipeline逻辑全程使用别名,和具体数据源完全解耦:

# src/<你的项目名>/pipelines/pl0/pipeline.py
from kedro.pipeline import Pipeline, node
from .nodes import process_pl0

def create_pipeline(**kwargs) -> Pipeline:
    return Pipeline(
        [
            node(
                func=process_pl0,
                inputs="pl0_raw_input",
                outputs="pl0_output_dataset",
                name="process_pl0_node",
            )
        ]
    )
# src/<你的项目名>/pipelines/pl1/pipeline.py
from kedro.pipeline import Pipeline, node
from .nodes import process_pl1

def create_pipeline(**kwargs) -> Pipeline:
    return Pipeline(
        [
            node(
                func=process_pl1,
                inputs="pl1_input_dataset", # 直接使用别名,不绑定具体数据源
                outputs="pl1_output_dataset",
                name="process_pl1_node",
            )
        ]
    )
# src/<你的项目名>/pipeline_registry.py
from kedro.pipeline import Pipeline
from .pipelines.pl0 import create_pipeline as create_pl0_pipeline
from .pipelines.pl1 import create_pipeline as create_pl1_pipeline

def register_pipelines() -> dict[str, Pipeline]:
    pl0_pipeline = create_pl0_pipeline()
    pl1_pipeline = create_pl1_pipeline()
    return {
        "pl0": pl0_pipeline,
        "pl1": pl1_pipeline,
        "__default__": pl0_pipeline + pl1_pipeline, # 串联两个Pipeline用于集成运行
    }

运行命令示例

  • 本地独立运行PL1:执行kedro run --pipeline pl1,默认加载local环境配置,自动读取测试CSV作为输入
  • 生产集成串联运行:执行kedro run --env prod,加载prod环境配置,PL1自动读取PL0的输出作为输入
  • 集成模式下单独运行PL1:执行kedro run --pipeline pl1 --env prod,直接读取已生成的PL0输出文件作为输入

注意事项

  • 使用内存数据集时必须全链路一次性运行,单独跑PL1会找不到输入,需要提前生成PL0的输出文件或者切换到CSV映射的配置
  • Kedro的多环境配置是叠加生效的,prod环境的同名配置会自动覆盖base中的定义,不需要重复配置公共字段

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 19:36:03