Kedro管道中如何更新数据集?规避循环与重复数据
解决Kedro中SQL表追加去重且避免循环依赖的方案
方案一:拆分数据集配置,用中间节点处理去重
核心思路是给同一张SQL表配置两个不同的数据集条目:一个只读,一个用于追加写入,避开Kedro的循环依赖限制。
修改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"编写节点逻辑:
分三步完成:读取已有日期、过滤重复数据、追加写入主表: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的「数据中心化」设计理念。
编写自定义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)配置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"简化节点逻辑:
节点只需生成待追加数据,直接输出到目标数据集即可: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在管道运行前查询已有日期,作为参数传入节点。
编写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)注册Hook:
在src/<project_name>/settings.py中添加:from .hooks import DataDedupHooks HOOKS = [DataDedupHooks()]节点使用参数:
节点通过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
相关产品推荐
相关产品推荐

