Σημείωμα
Η πρόσβαση σε αυτήν τη σελίδα απαιτεί εξουσιοδότηση. Μπορείτε να δοκιμάσετε να εισέλθετε ή να αλλάξετε καταλόγους.
Η πρόσβαση σε αυτήν τη σελίδα απαιτεί εξουσιοδότηση. Μπορείτε να δοκιμάσετε να αλλάξετε καταλόγους.
The create_auto_cdc_from_snapshot_flow function creates a flow that uses Lakeflow pipelines change data capture (CDC) functionality to process source data from database snapshots. See How AUTO CDC FROM SNAPSHOT works.
Note
This function replaces the previous function apply_changes_from_snapshot(). The two functions have the same signature. Databricks recommends updating to use the new name.
Important
You must have a target streaming table for this operation. To create the required target table, you can use the create_streaming_table() function.
Syntax
from pyspark import pipelines as dp
dp.create_auto_cdc_from_snapshot_flow(
target = "<target-table>",
source = Any,
keys = ["key1", "key2", "keyN"],
stored_as_scd_type = "1",
track_history_column_list = None,
track_history_except_column_list = None,
once = False
)
Note
For AUTO CDC FROM SNAPSHOT processing, the default behavior is to insert a new row when a matching record with the same key(s) does not exist in the target. If a matching record does exist, it is updated only if any of the values in the row have changed. Rows with keys present in the target but no longer present in the source are deleted.
To learn more about CDC processing with snapshots, see The AUTO CDC APIs: Simplify change data capture with pipelines. For examples of using the create_auto_cdc_from_snapshot_flow() function, see the periodic snapshot ingestion, historical snapshot ingestion, and SCD Type 1 backfill examples.
Parameters
| Parameter | Type | Description |
|---|---|---|
target |
str |
Required. The name of the table to be updated. You can use the create_streaming_table() function to create the target table before executing the create_auto_cdc_from_snapshot_flow() function. |
source |
str or lambda function |
Required. Either the name of a table or view to snapshot periodically or a Python lambda function that returns the snapshot DataFrame to be processed and the snapshot version. See Implement the source argument. |
keys |
list |
Required. The column or combination of columns that uniquely identify a row in the source data. This is used to identify which CDC events apply to specific records in the target table. You can specify either:
|
stored_as_scd_type |
str or int |
Whether to store records as SCD type 1 or SCD type 2. Set to 1 for SCD type 1 or 2 for SCD type 2. The default is SCD type 1. |
track_history_column_list or track_history_except_column_list |
list |
A subset of output columns to be tracked for history in the target table. Use track_history_column_list to specify the complete list of columns to be tracked. Use track_history_except_column_list to specify the columns to be excluded from tracking. You can declare either value as a list of strings or as Spark SQL col() functions:
Arguments to col() functions cannot include qualifiers. For example, you can use col(userId), but you cannot use col(source.userId). The default is to include all columns in the target table when no track_history_column_list or track_history_except_column_list argument is passed to the function. |
once |
bool |
Whether to run the snapshot flow only one time. When True, the flow is skipped during subsequent incremental updates after it commits successfully. A full refresh of the target re-runs the flow. Use this option for a one-time snapshot load or backfill. The default is False. |
Notes
For SCD Type 1 targets in triggered pipelines, one create_auto_cdc_from_snapshot_flow() flow can share a target with one or more AUTO CDC flows defined in Python or SQL. The following requirements apply:
- Give each
AUTO CDCflow a unique name. - Use the same number of keys in the same order for all flows. Snapshot-flow key names are compared case-insensitively with
AUTO CDCkey names. MultipleAUTO CDCflows must use identical key names and casing. - Use exactly the same data type for the snapshot version and the sequencing column of each
AUTO CDCflow. - Do not define expectations on the
AUTO CDC FROM SNAPSHOTflow. - Do not use
IGNORE NULL UPDATES. In Python, do not setignore_null_updates,ignore_null_updates_column_list, orignore_null_updates_except_column_liston thecreate_auto_cdc_flow()flows. - If you add an
AUTO CDCflow to an existingAUTO CDC FROM SNAPSHOTtarget, the target must not contain user columns whose names conflict with reservedAUTO CDCsystem columns. - Do not use this pattern for SCD Type 2 or bitemporal targets.
For an example, see Add a backfill to an AUTO CDC SCD Type 1 table.
Implement the source argument
The create_auto_cdc_from_snapshot_flow() function includes the source argument. For processing historical snapshots, the source argument is expected to be a Python lambda function that returns two values to the create_auto_cdc_from_snapshot_flow() function: a Python DataFrame containing the snapshot data to be processed and a snapshot version.
The first invocation of the function must return a snapshot DataFrame and version. If the first invocation returns None, the pipeline update fails. After at least one snapshot is processed, return None to indicate that no additional snapshots are available.
The following is the signature of the lambda function:
lambda Any => Optional[(DataFrame, Any)]
- The argument to the lambda function is the most recently processed snapshot version.
- The return value of the lambda function is
Noneor a tuple of two values: The first value of the tuple is a DataFrame containing the snapshot to be processed. The second value of the tuple is the snapshot version that represents the logical order of the snapshot.
An example that implements and calls the lambda function:
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Tuple[DataFrame, Optional[int]]:
if latest_snapshot_version is None:
return (spark.read.load("filename.csv"), 1)
else:
return None
create_auto_cdc_from_snapshot_flow(
# ...
source = next_snapshot_and_version,
# ...
)
The Lakeflow pipelines runtime performs the following steps each time the pipeline that contains the create_auto_cdc_from_snapshot_flow() function is triggered:
- Runs the
next_snapshot_and_versionfunction to load the next snapshot DataFrame and the corresponding snapshot version. - If the first invocation returns no DataFrame, the update fails. If a later invocation returns no DataFrame, the run terminates and the pipeline update is marked as complete.
- Detects the changes in the new snapshot and incrementally applies them to the target table.
- Returns to step #1 to load the next snapshot and its version.