szehon-ho commented on code in PR #57149: URL: https://github.com/apache/spark/pull/57149#discussion_r3752646276
########## docs/declarative-pipelines-programming-guide.md: ########## @@ -517,6 +517,360 @@ AS INSERT INTO customers_us SELECT * FROM STREAM(customers_us_east); ``` +## Change Data Capture (CDC) with Auto CDC + +Many source systems emit a stream of *change events* rather than a snapshot of the current data: each record describes an insert, update, or delete to a row, identified by a key. Applying these events correctly to a target table by hand is tricky. You have to match events to existing rows, apply them in the right order, and handle out-of-order and duplicate events without corrupting the table. + +**Auto CDC** does this for you. You point it at a source of change events and tell it how to identify and order them, and SDP maintains a target streaming table that always reflects the latest state for each key. + +### Keys and sequencing + +Auto CDC needs two things from you to make sense of a change feed: + +- The **keys** are the columns that identify a row across events. Events sharing a key describe the same logical row over time. In a customer feed, that is usually the customer id. +- The **sequencing expression** says what order the events for a key happened in. Its value for an event is that event's **sequence value**, and Auto CDC treats the highest sequence value it has seen for a key as the most recent state. Change feeds normally carry something suitable already: a monotonically increasing version or commit number, or a commit timestamp. + +Sequence values matter because a change feed does not have to arrive in order. Sequencing is what lets Auto CDC recognize that an event it just received is older than what it already applied, and place it correctly, rather than letting arrival order corrupt the table. + +### What Auto CDC does + +Given a stream of change events, Auto CDC keeps the target table in sync with the source: + +- **Inserts and updates** - For each key, the event with the highest sequence value wins. If no row exists for the key, it's inserted; if one exists, it's overwritten with the latest values. +- **Deletes** - Events that match a delete condition you supply remove the corresponding row from the target. +- **Out-of-order and duplicate events** - Events don't have to arrive in order, and you don't need to de-duplicate the source. An event whose sequence value is older than the state already applied for its key is discarded, and a re-delivered event converges to the same result. + +This behavior implements **Slowly Changing Dimensions (SCD) Type 1**: the target keeps only the current version of each row, with no history of prior values. SCD Type 1 is the only mode currently supported. + +For example, take these change events. Note that they are **not** in `version` order: the last event for `id = 1` is a stale update that arrives after the newer one. + +| id | name | version | op | +|----|--------|---------|--------| +| 1 | alice | 1 | UPSERT | +| 2 | bob | 1 | UPSERT | +| 1 | alicia | 2 | UPSERT | +| 2 | bob | 2 | DELETE | +| 3 | carol | 1 | UPSERT | +| 1 | alice | 1 | UPSERT | + +Auto CDC with `stored_as_scd_type=1`, keyed on `id` and sequenced by `version`, produces this target table: + +| id | name | version | +|----|--------|---------| +| 1 | alicia | 2 | +| 3 | carol | 1 | + +Walking through it by key: + +- **id 1** was inserted as `alice`, then updated to `alicia` at version 2. The re-delivered `alice` event at version 1 arrives last but is ignored, because version 1 is older than the version 2 already applied. The row keeps `alicia`. +- **id 2** was inserted, then deleted at version 2, so it is absent from the target. +- **id 3** was inserted and never changed. + +### Requirements + +- The **target must be a streaming table** that already exists in the pipeline. Create it with `create_streaming_table` (Python) or `CREATE STREAMING TABLE` (SQL) before defining the Auto CDC flow, or use the combined SQL form shown below that does both at once. +- The **target's format must support row-level operations.** Auto CDC maintains the target with MERGE, so the table must be backed by a connector implementing the DSv2 `SupportsRowLevelOperations` interface. A target that does not fails at startup with `AUTOCDC_TARGET_DOES_NOT_SUPPORT_MERGE`. Spark's built-in file formats, including Parquet, do **not** qualify; see [Choosing a target format](#choosing-a-target-format). +- The **source must be a streaming source** (read with `spark.readStream` in Python or `STREAM(...)` in SQL). CDC is an incremental operation over newly arriving change events. +- You must provide **keys** (one or more columns that identify a row) and a **sequencing expression** (used to order events per key). + +The sequencing expression may be any SQL expression over the source columns, not just a bare column reference, so `SEQUENCE BY` on a struct of `(commit_ts, seq_no)` or a cast is fine. It must satisfy three constraints: + +- **Its type must be orderable.** A non-orderable type fails with `AUTOCDC_MICROBATCH_VALIDATION.NON_ORDERABLE_SEQUENCE`. +- **It must never be null.** A microbatch containing a null sequence value fails with `AUTOCDC_MICROBATCH_VALIDATION.NULL_SEQUENCE` rather than guessing an order. +- **Its result type must stay the same across runs.** You may change the expression between incremental runs, but not its type, or recorded values would stop being comparable; that fails with `SEQUENCING_TYPE_DRIFT` and needs a full refresh. + +Ties are broken arbitrarily, so prefer an expression that is unique per event for a key if you need deterministic results across replays. + +### Defining an Auto CDC Flow in Python + +Use `create_auto_cdc_flow` to write change events into a target streaming table. Create the target with `create_streaming_table` first. + +```python +from pyspark import pipelines as dp + +# The source of change events: a streaming read of the CDC feed. [email protected] +def cdc_events(): + return spark.readStream.table("cdc_source") + +# The target that Auto CDC keeps in sync. It must be a streaming table, in a +# catalog whose format supports row-level operations (see "Choosing a target +# format" below). +dp.create_streaming_table("customers") + +# The Auto CDC flow that applies the change events to the target. +dp.create_auto_cdc_flow( + target="customers", + source="cdc_events", + keys=["id"], + sequence_by="version", + apply_as_deletes="op = 'DELETE'", + except_column_list=["op"], + stored_as_scd_type=1, +) +``` + +`create_auto_cdc_flow` accepts the following arguments: + +| Parameter | Required | Description | +|-----------|----------|-------------| +| `target` | Yes | Name of the target streaming table that receives the changes. It must already be defined in the pipeline. | +| `source` | Yes | Name of the CDC source dataset to stream change events from. | +| `keys` | Yes | The column or columns that uniquely identify a row. A list of column names (strings) or `Column` objects, given as unqualified identifiers: for example `"id"` or `col("id")`, but not `"source.id"`. | +| `sequence_by` | Yes | An expression used to order change events for each key. The highest value wins. A SQL expression string or a `Column`. | +| `apply_as_deletes` | No | A boolean expression identifying events that represent deletes. Matching rows are removed from the target. A SQL expression string or a `Column`. | +| `column_list` | No | The columns to include in the target. Mutually exclusive with `except_column_list`. | +| `except_column_list` | No | The columns to exclude from the target; all other columns are included. Mutually exclusive with `column_list`. Commonly used to drop operation/metadata columns such as `op`. | +| `stored_as_scd_type` | No | The SCD type of the target. Only `1` (or `"1"`) is supported. | +| `name` | No | The name of the flow. Defaults to the target table name. | +| `spark_conf` | No | Spark confs to set while the flow runs. These override confs set on the destination, the pipeline, or the cluster. | + +If you specify neither `column_list` nor `except_column_list`, all columns from the source are written to the target. That is usually not what you want for a CDC feed: the operation column and any other change-feed bookkeeping would land in the target alongside the data. Exclude them explicitly, as the examples here do with `except_column_list=["op"]`. + +`keys`, `sequence_by`, `column_list`, and `except_column_list` must be given as unqualified column identifiers: `"id"` or `col("id")`, but not `"cdc_events.id"` or `col("cdc_events.id")`. + +### Defining an Auto CDC Flow in SQL + +SQL provides two forms. The first attaches an Auto CDC flow to a streaming table you have already declared: + +```sql +CREATE STREAMING TABLE customers; + +CREATE FLOW customers_cdc AS AUTO CDC INTO customers +FROM STREAM(cdc_events) +KEYS (id) +APPLY AS DELETE WHEN op = 'DELETE' +SEQUENCE BY version +COLUMNS * EXCEPT (op); +``` + +The second declares the streaming table and its Auto CDC flow together: + +```sql +CREATE STREAMING TABLE customers +FLOW AUTO CDC +FROM STREAM(cdc_events) +KEYS (id) +APPLY AS DELETE WHEN op = 'DELETE' +SEQUENCE BY version +COLUMNS * EXCEPT (op); +``` + +`FROM STREAM(source)` and `KEYS (col, ...)` come first, in that order. The remaining clauses may appear in any order after them: + +- `FROM STREAM(source)` - the streaming CDC source. **Required, first.** +- `KEYS (col, ...)` - the key columns that identify a row. **Required, second.** +- `SEQUENCE BY expr` - the expression that orders events per key. **Required.** +- `APPLY AS DELETE WHEN condition` - marks events that represent deletes. Optional. +- `COLUMNS (col, ...)` or `COLUMNS * EXCEPT (col, ...)` - selects or excludes columns. Optional; if omitted, all source columns are written. +- `STORED AS SCD TYPE 1` - selects the SCD type. Optional; only Type 1 is supported. + +`CREATE FLOW ... AS AUTO CDC INTO` also accepts an optional `COMMENT`, and `CREATE STREAMING TABLE ... FLOW AUTO CDC` accepts `IF NOT EXISTS`. + +### Choosing a Target Format + +Auto CDC applies each microbatch to the target with a MERGE, so the target has to be a table that supports row-level updates and deletes. Concretely, its connector must implement the DSv2 `SupportsRowLevelOperations` interface. + +Spark's built-in file-based formats do not. Pointing an Auto CDC flow at a plain Parquet target - which is what a `create_streaming_table("customers")` with no `format` gives you - fails when the flow starts: + +``` +[AUTOCDC_TARGET_DOES_NOT_SUPPORT_MERGE] Cannot start AutoCDC flow: the target table +`spark_catalog`.`default`.`customers` (format: parquet) does not support row-level +operations. AutoCDC requires a target backed by a connector that supports MERGE. +``` + +To use Auto CDC you need a table provider that implements row-level operations, configured as a catalog in your pipeline. Lakehouse connectors such as Apache Iceberg are the usual choice; check your connector's documentation for whether it implements the DSv2 row-level operation interfaces and how to register its catalog. Note that supporting the `MERGE INTO` SQL statement is not on its own sufficient: a connector can implement MERGE through its own planner extension without implementing the DSv2 interface that Auto CDC requires. + +Configuring a catalog also gives you a persistent one, which incremental runs need. Spark's default session catalog keeps table metadata only for the life of the session, so a second `spark-pipelines run` cannot find the tables the first run created and fails with `LOCATION_ALREADY_EXISTS`. + +The examples below use a catalog named `lakehouse` to stand in for such a connector. Substitute your own catalog name and set the corresponding `spark.sql.catalog.*` configuration in your pipeline spec: + +```yaml +name: cdc_demo +storage: file:///absolute/path/to/storage/dir +catalog: lakehouse +database: cdc_demo +libraries: + - glob: + include: transformations/** +configuration: + spark.sql.catalog.lakehouse: <your connector's catalog class> +``` + +### End-to-End Example + +This example builds a small pipeline that ingests customer change events and maintains a `customers` table holding the latest state of each customer. Running it in two passes shows both that updates and deletes are applied incrementally, and that a late event is correctly ignored. + +It reads the change feed from a directory of JSON files, so you can append a batch and re-run to see what happens. The target uses the `lakehouse` catalog from [Choosing a target format](#choosing-a-target-format); substitute your own row-level-operation-capable catalog. + +Create a pipeline project: + +```bash +spark-pipelines init --name cdc_demo +cd cdc_demo +``` + +Point the pipeline at your catalog by editing `spark-pipeline.yml` as shown in [Choosing a target format](#choosing-a-target-format), then put the following in `transformations/customers_cdc.py`: + +```python +from pyspark import pipelines as dp +from pyspark.sql.types import ( + IntegerType, LongType, StringType, StructField, StructType) + +# An explicit schema keeps the streaming JSON read from having to infer one. +SCHEMA = StructType([ + StructField("id", IntegerType()), + StructField("name", StringType()), + StructField("version", LongType()), + StructField("op", StringType()), +]) + +# Ingest the raw change events. In a real pipeline this would read from Kafka, +# cloud storage, or a database CDC feed; here it tails a directory of JSON files. [email protected](name="cdc_events") +def cdc_events(): + return spark.readStream.schema(SCHEMA).json("file:///tmp/cdc_demo/events") + +# Declare the target streaming table that Auto CDC maintains. +dp.create_streaming_table("customers") + +# Apply the change events to the target. +dp.create_auto_cdc_flow( + target="customers", + source="cdc_events", + keys=["id"], + sequence_by="version", + apply_as_deletes="op = 'DELETE'", + except_column_list=["op"], + stored_as_scd_type=1, +) +``` + +Write the first batch of change events, inserting two customers: + +```bash +mkdir -p /tmp/cdc_demo/events +cat > /tmp/cdc_demo/events/batch1.json <<'EOF' +{"id": 1, "name": "alice", "version": 1, "op": "UPSERT"} +{"id": 2, "name": "bob", "version": 1, "op": "UPSERT"} +EOF +``` + +Run the pipeline: + +```bash +spark-pipelines run +``` + +`customers` now holds both rows, with the `op` column excluded: + +| id | name | version | +|----|-------|---------| +| 1 | alice | 1 | +| 2 | bob | 1 | + +Now add a second batch. It updates `id 1`, deletes `id 2`, inserts `id 3`, and ends with a **late duplicate** of the original `id 1` event, which is what a re-delivering source might send: + +```bash +cat > /tmp/cdc_demo/events/batch2.json <<'EOF' +{"id": 1, "name": "alicia", "version": 2, "op": "UPSERT"} +{"id": 2, "name": "bob", "version": 2, "op": "DELETE"} +{"id": 3, "name": "carol", "version": 1, "op": "UPSERT"} +{"id": 1, "name": "alice", "version": 1, "op": "UPSERT"} +EOF +``` + +Run the pipeline again. Because `customers` is a streaming table, this run processes only the new file: + +```bash +spark-pipelines run +``` + +| id | name | version | +|----|--------|---------| +| 1 | alicia | 2 | +| 3 | carol | 1 | + +`alice` became `alicia`, `bob` is gone, and `carol` was inserted. The late `alice` event at version 1 did not resurrect the old name: version 1 is below the version 2 already recorded for that key, so Auto CDC discarded it. Had the pipeline applied events in arrival order, `id 1` would have wrongly reverted to `alice`. + +### How-Tos + +#### Handling deletes + +Change feeds usually mark deletes with an operation column or a tombstone flag rather than removing the row. Give Auto CDC a boolean expression that identifies delete events with `apply_as_deletes` (Python) or `APPLY AS DELETE WHEN` (SQL): + +```python +dp.create_auto_cdc_flow( + target="customers", + source="cdc_events", + keys=["id"], + sequence_by="version", + apply_as_deletes="op = 'DELETE'", +) +``` + +When an event matches the delete condition, the row for its key is removed from the target. If you don't supply a delete condition, every event is treated as an insert or update. + +#### Selecting which columns land in the target + +CDC feeds often carry metadata columns (the operation type, a timestamp, source offsets) that you don't want in the target table. Use `except_column_list` / `COLUMNS * EXCEPT` to drop them, or `column_list` / `COLUMNS` to name exactly the columns to keep. The two options are mutually exclusive. + +```python +# Keep everything except the operation column. +dp.create_auto_cdc_flow( + target="customers", + source="cdc_events", + keys=["id"], + sequence_by="version", + except_column_list=["op"], +) + +# Or keep only an explicit set of columns. +dp.create_auto_cdc_flow( + target="customers", + source="cdc_events", + keys=["id"], + sequence_by="version", + column_list=["id", "name"], +) +``` + +#### Handling out-of-order and duplicate events + +You don't need to sort or de-duplicate the source. The target holds only the current state of each key, so a late event whose sequence value is older than what has already been applied for its key is discarded, and a re-delivered event converges to the same result. The end-to-end example above demonstrates both. Choose a `sequence_by` expression that strictly orders changes for a key, such as a monotonically increasing version number or a commit timestamp. + +#### Using a composite key + +Pass multiple columns to `keys` when a single column doesn't uniquely identify a row: + +```python +dp.create_auto_cdc_flow( + target="orders", + source="order_events", + keys=["region", "order_id"], + sequence_by="event_ts", +) +``` + +#### Changing the key set + +The set and types of `keys` are part of the flow's persisted state. Changing keys across incremental runs - renaming, swapping, adding, removing, or changing the type of a key column - is not supported and produces undefined results. To change the key set, [fully refresh](#spark-pipelines-run) the target table so it is recomputed from scratch: + +```bash +spark-pipelines run --full-refresh customers +``` + +### Auto CDC Considerations + +- **Target must be a streaming table** - You cannot apply an Auto CDC flow to a materialized view or an external table. Review Comment: Prefer avoid second person - line 865: "You cannot apply an Auto CDC flow to a materialized view or an external table" -> "An Auto CDC flow cannot target a materialized view or an external table" - line 870: "Your source cannot contain a column with that name" -> "The source cannot contain a column with that name" - line 871: "If you declare a schema on the target streaming table, it must currently include..." -> "A schema declared on the target streaming table must currently include..." -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
