如何在Kedro中增量更新管道输出的数据集而非覆盖?
如何在Kedro中增量更新管道输出的数据集而非覆盖?
看起来你遇到的问题是Kedro默认会覆盖输出数据集,而你需要把每次运行的新结果追加到result.csv里。我来给你一步步解决这个问题:
1. 调整数据集配置
首先,你需要让Kedro允许result.csv在不存在的时候也能正常加载(比如第一次运行的时候),修改你的catalog.yml里的{namespace}.result配置:
"{namespace}.result": type: pandas.CSVDataset filepath: data/01_raw/{namespace}/result.csv allow_missing: True # 新增这一行,允许文件不存在时返回None
2. 修改合并函数,加入增量逻辑
原来的concatenate_csvs只处理新生成的DataFrame,现在我们要把已有的result.csv数据也加进来合并。更新你的函数:
import pandas as pd def concatenate_csvs(existing_result, *new_dataframes): # 处理第一次运行,现有结果为空的情况 if existing_result is None or existing_result.empty: combined_df = pd.concat(new_dataframes, ignore_index=True) else: # 把现有数据和新数据合并 combined_df = pd.concat([existing_result, *new_dataframes], ignore_index=True) # 可选:如果担心重复数据,可以加个去重步骤 combined_df = combined_df.drop_duplicates(ignore_index=True) return combined_df
3. 调整管道定义,传递现有结果作为输入
现在你需要在管道里把已有的result数据集作为输入传给合并节点,然后把合并后的结果再输出回result。修改你的管道代码:
from kedro.pipeline import pipeline, node def create_pipeline(**kwargs): # 假设你已经有了各个子管道的输出节点,比如output1、output2... pipe_the_hacker_news = pipeline( [ # 其他节点... node( func=concatenate_csvs, # 第一个输入是现有result,后面是本次运行生成的所有新数据 inputs=["{namespace}.result", "output_of_pipeline1", "output_of_pipeline2"], outputs="{namespace}.result", name="concatenate_results" ) ] ) return pipe_the_hacker_news
原理说明
这样每次运行管道时,Kedro会先尝试读取现有的result.csv:
- 如果是第一次运行,文件不存在,
existing_result会是None,函数就直接合并本次的新数据生成第一个版本的result.csv - 后续运行时,会把现有数据和新数据合并,再覆盖写入
result.csv——但因为已经包含了之前的所有数据,所以相当于实现了增量更新
备注:内容来源于stack exchange,提问作者Lucas Queiroz
相关产品推荐
相关产品推荐

