CREAR FLUJO (canalizaciones)

Use la CREATE FLOW instrucción para crear flujos o rellenos para tablas en una canalización.

Syntax

CREATE FLOW flow_name [COMMENT comment] AS
{
  AUTO CDC [ONCE] INTO target_table create_auto_cdc_flow_spec |
  INSERT [ONCE] INTO target_table BY NAME [ replace_using_spec ] query
}

replace_using_spec
  REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column

Parámetros

  • flow_name

    Nombre del flujo que se va a crear.

  • COMENTARIO

    Descripción opcional del flujo.

  • CD AUTOMÁTICO EN

    Instrucción AUTO CDC ... INTO que define el flujo, con un create_auto_cdc_flow_spec. Debe incluir una instrucción AUTO CDC ... INTO o INSERT INTO. Use AUTO CDC ... INTO cuando la consulta de origen use la semántica de datos modificados.

    Para más información, consulte AUTO CDC INTO (canalizaciones).

  • target_table

    Tabla que se va a actualizar. Debe ser una tabla de Streaming.

  • INSERT EN

    Define una consulta de tabla que se inserta en la tabla de destino. Si no se proporciona la ONCE opción , la consulta debe ser una consulta de streaming . Use la palabra clave STREAM para usar la semántica de streaming para leer desde el origen. Si la lectura encuentra un cambio o eliminación en un registro existente, se produce un error. Es más seguro leer de orígenes estáticos o de solo anexión. Para ingerir datos que tienen confirmaciones de cambios, puede usar Python y la skipChangeCommits opción para controlar errores.

    INSERT INTO es mutuamente excluyente con AUTO CDC ... INTO. Use AUTO CDC ... INTO cuando los datos de origen incluyan la funcionalidad de captura de datos modificados (CDC). Utilice INSERT INTO cuando el origen no lo haga.

    Para más información sobre los datos de streaming, consulte Transformación de datos con canalizaciones.

  • REEMPLAZAR USANDO ( column_name [, ...] ) SECUENCIA POR sequence_column

    Importante

    Esta característica se encuentra en su versión beta. Requiere Databricks en tiempo de ejecución 18.2 y superiores.

    Define el flujo como un REPLACE USING flujo, que reemplaza todas las filas de la tabla objetivo que coincidan con las columnas clave especificadas y deja todas las demás filas intactas. Úsalo REPLACE USING cuando tu fuente es una serie de instantáneas parciales codificadas por columna. SEQUENCE BY ordena las actualizaciones para que gane la secuencia más alta para una clave, incluso cuando las actualizaciones llegan fuera de orden.

    Especifica al menos una columna clave y exactamente una SEQUENCE BY columna. La consulta debe ser una consulta de streaming y BY NAME es obligatoria. REPLACE USING no puede combinarse con ONCE ni con AUTO CDC ... INTO.

    Para más información, véase Sustitución parcial de instantáneas con REEMPLAZAR flujos UTILIZADOS.

  • Una vez

    Opcionalmente, defina el flujo como un flujo de una sola vez, como un reposición. El uso de ONCE cambia el flujo de dos maneras:

    • El origen query o create_auto_cdc_flow_spec no es una tabla de streaming.
    • El flujo se ejecuta una vez de forma predeterminada. Si la canalización se actualiza por completo, el flujo ONCE se ejecuta nuevamente para recrear los datos.

    ONCE No se puede usar con REPLACE USING, que requiere una fuente de streaming.

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);