如何为Kedro Catalog动态传递save_args以实现Delta表的增量分区写入
replaceWhere for Delta Tables in Kedro Great question! The core issue here is that Kedro's Catalog is statically initialized when your pipeline starts, but your replaceWhere filter relies on runtime values generated by the meta_reload node. You can't directly reference that dynamic date in your YAML config, but there are two straightforward ways to handle this:
Approach 1: Manually Save the Dataset in Your Node
This is the most direct method—you'll bypass Kedro's automatic output binding and explicitly handle the save operation in your node function, using the runtime date from meta_reload to build the replaceWhere clause.
Step 1: Update Your Catalog Config
First, simplify your my_dataset config by removing the hardcoded replaceWhere (we'll set this dynamically later):
my_dataset: type: my_project.io.pyspark.SparkDataSet filepath: "s3://${bucket_de_pipeline}/${data_environment_project}/${data_environment_intermediate}/my_dataset/" file_format: delta layer: intermediate save_args: mode: "overwrite" partitionBy: [ "DATE_ID" ]
Step 2: Modify Your Node Function
In your node, fetch the start date from meta_reload, filter your data, then manually load the dataset and save it with dynamic save_args:
from kedro.io import DataCatalog from pyspark.sql import DataFrame def process_incremental_data( input_data: DataFrame, meta_reload_data: DataFrame, data_catalog: DataCatalog ): # Extract the start date from your meta_reload dataset # Adjust this based on how your meta data is structured start_date_row = meta_reload_data.select("start_date").collect()[0] start_date = start_date_row["start_date"] # Filter your working dataset using the dynamic start date filtered_data = input_data.filter(f"DATE_ID > '{start_date}'") # Load the target Delta dataset from the catalog delta_dataset = data_catalog.load("my_dataset") # Build dynamic save_args with the replaceWhere clause dynamic_save_args = { **delta_dataset.save_args, # Inherit existing save_args from catalog "replaceWhere": f"DATE_ID > '{start_date}'" } # Save the filtered data with the custom save_args delta_dataset.save(filtered_data, save_args=dynamic_save_args)
Step 3: Update Your Pipeline Definition
Since you're handling the save manually, your node doesn't need to declare my_dataset as an output. Your pipeline might look like this:
from kedro.pipeline import Pipeline, node def create_pipeline(**kwargs) -> Pipeline: return Pipeline( [ node( func=process_incremental_data, inputs=["input_working_dataset", "meta_reload", "data_catalog"], outputs=None, # No output binding—we save manually name="process_and_save_delta" ) ] )
Approach 2: Use a Kedro Hook to Modify Save Args at Runtime
If you need this dynamic behavior across multiple Delta datasets, a hook is a cleaner solution. Hooks let you intercept dataset operations and modify their configs before saving.
Step 1: Define a Custom Hook
Create a hook that runs before a dataset is saved, fetches the meta_reload date, and updates the replaceWhere clause:
from kedro.framework.hooks import hook_impl from kedro.io import DataCatalog from my_project.io.pyspark.SparkDataSet import SparkDataSet class DeltaDynamicSaveHook: @hook_impl def before_dataset_saved(self, dataset_name: str, dataset: SparkDataSet, data, catalog: DataCatalog): # Target only your specific Delta dataset if dataset_name == "my_dataset" and dataset.file_format == "delta": # Load the meta_reload data to get the start date meta_data = catalog.load("meta_reload") start_date = meta_data.select("start_date").collect()[0]["start_date"] # Update the save_args with dynamic replaceWhere dataset.save_args["replaceWhere"] = f"DATE_ID > '{start_date}'"
Step 2: Register the Hook
Add your hook to your project's hooks.py file so Kedro loads it at startup:
from .delta_hooks import DeltaDynamicSaveHook hooks = [DeltaDynamicSaveHook()]
Key Notes
- Why Catalog Can't Do This Natively: Kedro parses the Catalog YAML at pipeline initialization, before any nodes run. Since
meta_reloadgenerates the date at runtime, it's unavailable when the Catalog is loaded. - Spark Date Formatting: Ensure your
start_datematches the format of yourDATE_IDcolumn (e.g., ifDATE_IDis a string like20240520, formatstart_dateaccordingly).
内容的提问来源于stack exchange,提问作者Sandeep Gunda

