Note
Access to this page requires authorization. You can try signing in or changing directories.
Access to this page requires authorization. You can try changing directories.
Use the CREATE FLOW statement to create flows or backfills for tables in a pipeline.
Syntax
CREATE FLOW flow_name [COMMENT comment] AS
{
AUTO CDC [ONCE] INTO target_table create_auto_cdc_flow_spec |
AUTO CDC [ONCE] INTO target_table create_auto_cdc_from_snapshot_spec |
INSERT [ONCE] INTO target_table BY NAME [ replace_using_spec ] query
}
create_auto_cdc_from_snapshot_spec
FROM SNAPSHOT ( snapshot_query )
[ WITH VERSION ( version_query ) ]
KEYS ( key [, ...] )
[ STORED AS { SCD TYPE 1 | SCD TYPE 2 } ]
[ TRACK HISTORY ON { col_list | * EXCEPT ( col_list ) } ]
replace_using_spec
REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column
Parameters
flow_name
The name of the flow to create.
COMMENT
An optional description for the flow.
-
An
AUTO CDC ... INTOstatement that defines the flow, with acreate_auto_cdc_flow_spec. You must either include anAUTO CDC ... INTOstatement, or anINSERT INTOstatement. UseAUTO CDC ... INTOwhen the source query uses change data semantics.For more information, see AUTO CDC INTO (pipelines).
AUTO CDC ... FROM SNAPSHOT
An
AUTO CDC ... INTOstatement that derives changes by comparing snapshots instead of reading a change feed. Use this form when change data capture is not enabled on the source and only full snapshots are available. The source is specified in two parts: a requiredFROM SNAPSHOT (snapshot_query)clause that reads the snapshot data, and an optionalWITH VERSION (version_query)clause that selects the next snapshot version to process. See How AUTO CDC FROM SNAPSHOT works.FROM SNAPSHOT (snapshot_query)
Required. A query that reads the snapshot data for the version selected by
WITH VERSION (...). The engine diffs the result against the previously committed snapshot to derive inserts, updates, and deletes, and merges them into the target usingKEYSfor row identity andSTORED ASto determine how changes are stored.Call current_snapshot_version() inside this query to reference the version selected by
WITH VERSION (...). IfWITH VERSION (...)is not specified,current_snapshot_version()is not callable insideFROM SNAPSHOT (...).When
WITH VERSION (...)is omitted, the engine reads the source directly throughFROM SNAPSHOT (...), and the snapshot query runs only during the initial load, while the target has no committed data and no committed snapshot state. On any later update, when the target already contains data or has committed snapshot state, the flow fails withAUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. To process snapshots across multiple updates, useWITH VERSION (...).WITH VERSION (version_query)
Optional. A query that selects the next snapshot version to process. It must return exactly one column of an orderable type and either 0 or 1 rows. When it returns 1 row, the value must be non-null. The column can be a scalar value, such as a
BIGINT, or aSTRUCTwhose fields are all orderable. A version query that returns more than one column, more than one row, or a null value fails the flow withINVALID_AUTO_CDC_FROM_SNAPSHOT_VERSION_QUERY.During one pipeline update, the engine repeats the following steps: it evaluates the version query; if the query returns 0 rows, it stops processing this flow for the current update; if the query returns 1 row, the engine exposes that value through current_snapshot_version(), evaluates the snapshot query, commits the resulting snapshot, and exposes the committed version through last_snapshot_version(). The engine then re-evaluates the version query to select the next version. A single pipeline update processes versions in order until the version query returns no rows.
Every version returned after a successful commit must be greater than the previously committed version; a non-increasing version fails the update with
APPLY_CHANGES_FROM_SNAPSHOT_ERROR.OUT_OF_ORDER_SNAPSHOT_VERSION. The version value's data type must remain unchanged across snapshot commits; a data type change fails the update withAUTO_CDC_FROM_SNAPSHOT_VERSION_SCHEMA_CHANGED. A full refresh clears the persisted version state.KEYS
Required. The primary key columns used to identify rows across snapshots for change detection.
STORED AS { SCD TYPE 1 | SCD TYPE 2 }Optional. Specifies how changes are stored in the target table. The default is
SCD TYPE 1.TRACK HISTORY ON { col_list | * EXCEPT (col_list) }Optional. Applies only with
SCD TYPE 2. Specifies which columns trigger a new history row when they change. Provide either an explicit column list or* EXCEPT (col_list)to track every column except the ones listed.
Snapshot CDC does not support
WHEREorSEQUENCE BY. Cross-snapshot ordering is expressed throughWITH VERSION (...).target_table
The table to update. This must be a Streaming table.
INSERT INTO
Defines a table query that is inserted into to the target table. If the
ONCEoption is not supplied, the query must be a streaming query. Use the STREAM keyword to use streaming semantics to read from the source. If the read encounters a change or deletion to an existing record, an error is thrown. It is safest to read from static or append-only sources. To ingest data that has change commits, you can use Python and theskipChangeCommitsoption to handle errors.INSERT INTOis mutually exclusive withAUTO CDC ... INTO. UseAUTO CDC ... INTOwhen the source data includes change data capture (CDC) functionality. UseINSERT INTOwhen the source does not.For more information on streaming data, see Transform data with pipelines.
REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column
Important
This feature is in Beta. Requires Databricks Runtime 18.2 and above.
Defines the flow as a
REPLACE USINGflow, which replaces all rows in the target table matching the specified key columns and leaves all other rows untouched. UseREPLACE USINGwhen your source is a series of partial snapshots keyed by column.SEQUENCE BYorders the updates so the highest sequence for a key wins, even when updates arrive out of order.Specify at least one key column and exactly one
SEQUENCE BYcolumn. The query must be a streaming query, andBY NAMEis required.REPLACE USINGcan't be combined withONCEor withAUTO CDC ... INTO.For more information, see Partial snapshot replacement with REPLACE USING flows.
ONCE
Optionally define the flow as a one time flow, such as a backfill. Using
ONCEchanges the flow in two ways:- The source
queryorcreate_auto_cdc_flow_specis not a streaming table. - The flow is run one time by default. If the pipeline is updated with a complete refresh, then the
ONCEflow runs again to recreate the data.
ONCEcan't be used withREPLACE USING, which requires a streaming source.- The source
Examples
-- EXAMPLE 1:
-- Create a streaming table, and add two flows that append data to it:
CREATE OR REFRESH STREAMING TABLE users;
-- first flow into target_table:
CREATE FLOW users_flow AS
INSERT INTO users BY NAME
SELECT * FROM stream(raw_data.users);
-- second flow into target_table:
CREATE FLOW backfill_users AS
INSERT ONCE INTO users BY NAME
SELECT * FROM user_backfill_table;
-- EXAMPLE 2:
-- Create a streaming table, and add a flow that applies CDC changes to it:
CREATE OR REFRESH STREAMING TABLE admins_cdc_target_table;
-- first flow into target_table:
CREATE FLOW admin_cdc_flow AS
AUTO CDC INTO admins_cdc_target_table
FROM stream(cdc_data.admins)
KEYS (userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;
-- EXAMPLE 3:
-- Create a streaming table, and add a REPLACE USING flow that keeps the latest
-- row for each payment_id from a stream of partial snapshots:
CREATE OR REFRESH STREAMING TABLE payments_latest;
CREATE FLOW payments_replace_flow AS
INSERT INTO payments_latest BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
-- EXAMPLE 4:
-- AUTO CDC FROM SNAPSHOT without WITH VERSION: a one-time initial load from a snapshot table.
-- To process later snapshots on each update, add WITH VERSION (see EXAMPLE 5).
CREATE STREAMING TABLE users (user_id INT, name STRING, email STRING);
CREATE FLOW users_snapshot_flow AS
AUTO CDC ONCE INTO users
FROM SNAPSHOT (SELECT * FROM catalog.schema.users_snapshot)
KEYS (user_id)
STORED AS SCD TYPE 1;
-- EXAMPLE 5:
-- AUTO CDC FROM SNAPSHOT with WITH VERSION: pick the next file, then read it as the snapshot:
CREATE STREAMING TABLE orders (order_id INT, product STRING, quantity INT, order_date DATE);
CREATE FLOW orders_cdc AS
AUTO CDC INTO orders
FROM SNAPSHOT (
SELECT order_id, product, quantity, order_date
FROM read_files('/Volumes/catalog/schema/landing/orders/', format => 'json')
WHERE _metadata.file_path = (SELECT version.path FROM current_snapshot_version())
)
WITH VERSION (
SELECT struct(modification_time, path) AS version
FROM list_files('/Volumes/catalog/schema/landing/orders/')
WHERE (
NOT EXISTS (SELECT 1 FROM last_snapshot_version())
OR struct(modification_time, path) > (SELECT version FROM last_snapshot_version())
)
ORDER BY modification_time, path
LIMIT 1
)
KEYS (order_id)
STORED AS SCD TYPE 2;
-- EXAMPLE 6:
-- One-time snapshot backfill plus a streaming CDC flow into the same target.
-- The backfill omits WITH VERSION, so it uses an implicit timestamp version. The
-- streaming flow's SEQUENCE BY column (event_ts) must be a TIMESTAMP so its type
-- matches that implicit version on the shared target.
CREATE STREAMING TABLE customers (
customer_id INT, name STRING, email STRING, address STRING, event_ts TIMESTAMP
);
CREATE FLOW customers_snapshot_backfill AS
AUTO CDC ONCE INTO customers
FROM SNAPSHOT (SELECT * FROM catalog.schema.customers_snapshot)
KEYS (customer_id)
STORED AS SCD TYPE 1;
CREATE FLOW customers_cdc AS
AUTO CDC INTO customers
FROM STREAM(customers_cdc_events)
KEYS (customer_id)
SEQUENCE BY event_ts
STORED AS SCD TYPE 1;
For more about combining a one-time backfill with ongoing CDC on the same target, see Backfilling historical data with pipelines.