CREATE FLOW (pipelines)

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.

  • AUTO CDC INTO

    An AUTO CDC ... INTO statement that defines the flow, with a create_auto_cdc_flow_spec. You must either include an AUTO CDC ... INTO statement, or an INSERT INTO statement. Use AUTO CDC ... INTO when the source query uses change data semantics.

    For more information, see AUTO CDC INTO (pipelines).

  • AUTO CDC ... FROM SNAPSHOT

    An AUTO CDC ... INTO statement 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 required FROM SNAPSHOT (snapshot_query) clause that reads the snapshot data, and an optional WITH 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 using KEYS for row identity and STORED AS to determine how changes are stored.

      Call current_snapshot_version() inside this query to reference the version selected by WITH VERSION (...). If WITH VERSION (...) is not specified, current_snapshot_version() is not callable inside FROM SNAPSHOT (...).

      When WITH VERSION (...) is omitted, the engine reads the source directly through FROM 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 with AUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. To process snapshots across multiple updates, use WITH 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 a STRUCT whose fields are all orderable. A version query that returns more than one column, more than one row, or a null value fails the flow with INVALID_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 with AUTO_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 WHERE or SEQUENCE BY. Cross-snapshot ordering is expressed through WITH 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 ONCE option 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 the skipChangeCommits option to handle errors.

    INSERT INTO is mutually exclusive with AUTO CDC ... INTO. Use AUTO CDC ... INTO when the source data includes change data capture (CDC) functionality. Use INSERT INTO when 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 USING flow, which replaces all rows in the target table matching the specified key columns and leaves all other rows untouched. Use REPLACE USING when your source is a series of partial snapshots keyed by column. SEQUENCE BY orders 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 BY column. The query must be a streaming query, and BY NAME is required. REPLACE USING can't be combined with ONCE or with AUTO 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 ONCE changes the flow in two ways:

    • The source query or create_auto_cdc_flow_spec is not a streaming table.
    • The flow is run one time by default. If the pipeline is updated with a complete refresh, then the ONCE flow runs again to recreate the data.

    ONCE can't be used with REPLACE USING, which requires a streaming 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.