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

Kedro管道中如何更新数据集?规避循环与重复数据

解决Kedro中SQL表追加去重且避免循环依赖的方案

方案一:拆分数据集配置,用中间节点处理去重

核心思路是给同一张SQL表配置两个不同的数据集条目:一个只读,一个用于追加写入,避开Kedro的循环依赖限制。

  1. 修改catalog.yml:
    新增只读数据集指向目标表,保留原数据集的append模式:

    my_main_dataset_read:
      type: pandas.SQLTableDataSet
      table_name: your_table_name
      credentials: db_credentials
      load_args:
        columns: ["date_column"]  # 仅加载日期列,减少数据传输量
    
    my_main_dataset_append:
      type: pandas.SQLTableDataSet
      table_name: your_table_name
      credentials: db_credentials
      save_args:
        if_exists: "append"
    
  2. 编写节点逻辑:
    分三步完成:读取已有日期、过滤重复数据、追加写入主表:

    def get_existing_dates(existing_data: pd.DataFrame) -> set:
        return set(existing_data["date_column"].dt.date)
    
    def filter_duplicate_dates(raw_new_data: pd.DataFrame, existing_dates: set) -> pd.DataFrame:
        return raw_new_data[~raw_new_data["date_column"].dt.date.isin(existing_dates)]
    
    def create_pipeline(**kwargs) -> Pipeline:
        return Pipeline(
            [
                node(
                    func=get_existing_dates,
                    inputs="my_main_dataset_read",
                    outputs="existing_dates",
                    name="get_existing_dates_node",
                ),
                node(
                    func=filter_duplicate_dates,
                    inputs=["raw_new_data", "existing_dates"],
                    outputs="filtered_new_data",
                    name="filter_duplicate_dates_node",
                ),
                node(
                    func=lambda df: df,
                    inputs="filtered_new_data",
                    outputs="my_main_dataset_append",
                    name="append_to_main_table_node",
                ),
            ]
        )
    

    整个DAG是线性无环的,完全符合Kedro的要求。

方案二:自定义数据集,封装去重逻辑到数据层

把去重逻辑直接封装在DataSet里,节点只需输出数据,自动完成去重追加,更贴合Kedro的「数据中心化」设计理念。

  1. 编写自定义DataSet:
    在src/<project_name>/extras/datasets/sql_table_with_dedup.py中创建:

    from typing import Any, Dict, Optional
    import pandas as pd
    from kedro.io import SQLTableDataSet
    
    class SQLTableWithDedupDataSet(SQLTableDataSet):
        def __init__(
            self,
            table_name: str,
            credentials: Dict[str, Any],
            dedup_column: str,
            load_args: Optional[Dict[str, Any]] = None,
            save_args: Optional[Dict[str, Any]] = None,
        ):
            super().__init__(
                table_name=table_name,
                credentials=credentials,
                load_args=load_args or {},
                save_args=save_args or {"if_exists": "append"},
            )
            self.dedup_column = dedup_column
    
        def _save(self, data: pd.DataFrame) -> None:
            # 查询已有去重列数据
            existing_query = f"SELECT DISTINCT {self.dedup_column} FROM {self.table_name}"
            with self._get_sql_engine().connect() as conn:
                existing_data = pd.read_sql(existing_query, conn)
            # 过滤重复数据
            filtered_data = data[~data[self.dedup_column].isin(existing_data[self.dedup_column])]
            # 非空时才追加写入
            if not filtered_data.empty:
                super()._save(filtered_data)
    
  2. 配置catalog.yml:
    替换原数据集为自定义类型:

    my_main_dataset:
      type: <project_name>.extras.datasets.sql_table_with_dedup.SQLTableWithDedupDataSet
      table_name: your_table_name
      credentials: db_credentials
      dedup_column: "date_column"  # 指定去重用的日期列
      save_args:
        if_exists: "append"
    
  3. 简化节点逻辑:
    节点只需生成待追加数据,直接输出到目标数据集即可:

    def generate_new_data() -> pd.DataFrame:
        # 生成或获取待追加数据的逻辑
        return new_data_df
    
    def create_pipeline(**kwargs) -> Pipeline:
        return Pipeline(
            [
                node(
                    func=generate_new_data,
                    inputs=None,
                    outputs="my_main_dataset",
                    name="generate_and_append_data_node",
                ),
            ]
        )
    

方案三:通过Hooks传递已有日期参数

如果主表数据量极大,不想读取全量数据,可以用Kedro Hooks在管道运行前查询已有日期,作为参数传入节点。

  1. 编写Hook:
    在src/<project_name>/hooks.py中添加:

    from kedro.framework.hooks import hook_impl
    from sqlalchemy import create_engine
    import pandas as pd
    
    class DataDedupHooks:
        @hook_impl
        def before_pipeline_run(self, run_params, pipeline, catalog):
            # 从catalog获取数据库凭证
            db_creds = catalog._data_sets["my_main_dataset"]._credentials
            engine = create_engine(f"{db_creds['dialect']}+{db_creds['driver']}://{db_creds['user']}:{db_creds['password']}@{db_creds['host']}:{db_creds['port']}/{db_creds['database']}")
            # 查询已有日期
            with engine.connect() as conn:
                existing_dates = pd.read_sql("SELECT DISTINCT date_column FROM your_table_name", conn)
            # 将日期集合存入run_params传递给节点
            run_params["existing_dates"] = set(existing_dates["date_column"].dt.date)
    
  2. 注册Hook:
    在src/<project_name>/settings.py中添加:

    from .hooks import DataDedupHooks
    
    HOOKS = [DataDedupHooks()]
    
  3. 节点使用参数:
    节点通过run_params获取已有日期,过滤新数据:

    def generate_and_filter_data(run_params: Dict) -> pd.DataFrame:
        existing_dates = run_params["existing_dates"]
        # 生成新数据
        new_data_df = ...
        # 过滤重复日期
        return new_data_df[~new_data_df["date_column"].dt.date.isin(existing_dates)]
    
    def create_pipeline(**kwargs) -> Pipeline:
        return Pipeline(
            [
                node(
                    func=generate_and_filter_data,
                    inputs="params:run_params",
                    outputs="my_main_dataset",
                    name="generate_filtered_data_node",
                ),
            ]
        )
    

内容的提问来源于stack exchange,提问作者Jean-Baptiste

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 20:05:39