Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Use la CREATE FLOW instrucción para crear flujos o rellenos para tablas en una canalización.
Note
CREATE FLOW Dirigirse a una tabla de streaming soporta tanto AUTO CDC ... INTO flujos y REPLACE WHERE S. Una tabla gestionada creada con CREATE TABLE ... FLOW no soporta captura de datos de cambio: un AUTO CDC ... INTO flujo contra una tabla gestionada falla con MANAGED_TABLE_DOES_NOT_SUPPORT_CDC.
CREATE FLOW Into a Streaming Table no está sujeto a esa limitación. Véase CREATE TABLE ... FLOW (tuberías).
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 | REPLACE WHERE condition ] 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
Parámetros
flow_name
Nombre del flujo que se va a crear.
COMENTARIO
Descripción opcional del flujo.
-
Instrucción
AUTO CDC ... INTOque define el flujo, con uncreate_auto_cdc_flow_spec. Debe incluir una instrucciónAUTO CDC ... INTOoINSERT INTO. UseAUTO CDC ... INTOcuando la consulta de origen use la semántica de datos modificados.Para más información, consulte AUTO CDC INTO (canalizaciones).
AUTO CDC ... DE INSTANTÁNEA
Una
AUTO CDC ... INTOsentencia que deriva cambios comparando instantáneas en lugar de leer un feed de cambios. Utiliza este formulario cuando la captura de datos de cambios no esté habilitada en la fuente y solo haya instantáneas completas disponibles. La fuente se especifica en dos partes: una cláusula obligatoriaFROM SNAPSHOT (snapshot_query)que lee los datos de la instantánea y una cláusula opcionalWITH VERSION (version_query)que selecciona la siguiente versión de la instantánea a procesar. Consulte Funcionamiento de AUTO CDC FROM SNAPSHOT.DE LA INSTANTÁNEA (snapshot_query)
Required. Una consulta que lee los datos de instantánea para la versión seleccionada por
WITH VERSION (...). El motor diferencia el resultado con la instantánea previamente comprometida para derivar insertos, actualizaciones y eliminaciones, y los fusiona en el destino usandoKEYSla identidad de fila ySTORED ASpara determinar cómo se almacenan los cambios.Llama a current_snapshot_version() dentro de esta consulta para referenciar la versión seleccionada por
WITH VERSION (...). SiWITH VERSION (...)no se especifica,current_snapshot_version()no se puede llamar dentroFROM SNAPSHOT (...)de .Cuando
WITH VERSION (...)se omite, el motor lee la fuente directamente a través deFROM SNAPSHOT (...), y la consulta de instantánea solo se ejecuta durante la carga inicial, mientras que el destino no tiene datos comprometidos ni estado de instantánea comprometido. En cualquier actualización posterior, cuando el destino ya contiene datos o tiene estado de instantánea comprometido, el flujo falla conAUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. Para procesar instantáneas a través de múltiples actualizaciones, utilizaWITH VERSION (...).CON VERSIÓN (version_query)
Optional. Una consulta que selecciona la siguiente versión snapshot para procesar. Debe devolver exactamente una columna de tipo ordenable y 0 o 1 fila. Cuando devuelve 1 fila, el valor debe ser no nulo. La columna puede ser un valor escalar, como un
BIGINT, o unSTRUCTcuyos campos son todos ordenables. Una consulta de versión que devuelve más de una columna, más de una fila o un valor nulo falla el flujo conINVALID_AUTO_CDC_FROM_SNAPSHOT_VERSION_QUERY.Durante una actualización de la canalización, el motor repite los siguientes pasos: evalúa la consulta de versión; si la consulta devuelve 0 filas, deja de procesar este flujo para la actualización actual; si la consulta devuelve 1 fila, el motor expone ese valor a través de current_snapshot_version(), evalúa la consulta de instantánea, compromete la instantánea resultante y expone la versión comprometida a través de last_snapshot_version(). El motor luego reevalúa la consulta de versión para seleccionar la siguiente versión. Una única actualización de pipeline procesa las versiones en orden hasta que la consulta de versiones no devuelve filas.
Cada versión devuelta tras un commit exitoso debe ser mayor que la versión previamente confirmada; una versión no creciente falla la actualización con
APPLY_CHANGES_FROM_SNAPSHOT_ERROR.OUT_OF_ORDER_SNAPSHOT_VERSION. El tipo de dato del valor de la versión debe permanecer sin cambios entre los commits de instantánea; un cambio de tipo de dato falla la actualización conAUTO_CDC_FROM_SNAPSHOT_VERSION_SCHEMA_CHANGED. Una actualización completa elimina el estado persistente de la versión.LLAVES
Required. Las columnas clave principales se usan para identificar filas a través de instantáneas para la detección de cambios.
STORED AS { SCD TYPE 1 | SCD TYPE 2 }Optional. Especifica cómo se almacenan los cambios en la tabla de destino. El valor predeterminado es
SCD TYPE 1.TRACK HISTORY ON { col_list | * EXCEPT (col_list) }Optional. Se aplica solo con
SCD TYPE 2. Especifica qué columnas activan una nueva fila de historial cuando cambian. Proporciona una lista explícita de columnas o* EXCEPT (col_list)para hacer seguimiento de todas las columnas excepto las que se mencionan.
Snapshot CDC no soporta
WHEREniSEQUENCE BY. El orden cruzado de instantáneas se expresa medianteWITH VERSION (...).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
ONCEopció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 laskipChangeCommitsopción para controlar errores.INSERT INTOes mutuamente excluyente conAUTO CDC ... INTO. UseAUTO CDC ... INTOcuando los datos de origen incluyan la funcionalidad de captura de datos modificados (CDC). UtiliceINSERT INTOcuando 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 USINGflujo, que reemplaza todas las filas de la tabla objetivo que coincidan con las columnas clave especificadas y deja todas las demás filas intactas. ÚsaloREPLACE USINGcuando tu fuente es una serie de instantáneas parciales codificadas por columna.SEQUENCE BYordena 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 BYcolumna. La consulta debe ser una consulta de streaming yBY NAMEes obligatoria.REPLACE USINGno puede combinarse conONCEni conAUTO CDC ... INTO.Para más información, véase Sustitución parcial de instantáneas con REEMPLAZAR flujos UTILIZADOS.
Condición REPLACE WHERE
Define el flujo como un
REPLACE WHEREflujo, que recalcula y sobrescribe un subconjunto objetivo de la tabla objetivo. En cada actualización, todas las filas de la tabla de coincidenciaconditionobjetivo se eliminan, la consulta de origen se recalcula para ese mismo rango de predicados y se insertan los resultados. Las filas que no coincidenconditionquedan intactas. No necesitas añadir el predicado a la consulta fuente; el motor de pipeline lo aplica automáticamente al leer desde la fuente.REPLACE WHEREUtiliza semántica por lotes, por lo que la consulta de código fuente no tiene que ser una consulta de streaming.BY NAMEes obligatorio.REPLACE WHERENo se puede combinar conONCE,REPLACE USINGoAUTO CDC ... INTO.Para más información, consulte procesamiento por lotes con flujos REPLACE WHERE y flujos REPLACE WHERE para tablas de streaming independientes.
Una vez
Opcionalmente, defina el flujo como un flujo de una sola vez, como un reposición. El uso de
ONCEcambia el flujo de dos maneras:- El origen
queryocreate_auto_cdc_flow_specno es una tabla de streaming. - El flujo se ejecuta una vez de forma predeterminada. Si la canalización se actualiza por completo, el flujo
ONCEse ejecuta nuevamente para recrear los datos.
ONCENo se puede usar conREPLACE USING, que requiere una fuente de streaming.- El origen
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;
-- EXAMPLE 7:
-- Create a streaming table, then add a REPLACE WHERE flow that recomputes and
-- overwrites a targeted window of the target table on each update:
CREATE STREAMING TABLE payments_latest;
CREATE FLOW payments_latest AS
INSERT INTO payments_latest BY NAME
REPLACE WHERE payment_date >= date_add(current_date(), -7)
SELECT payment_id, booking_id, status, payment_date
FROM samples.wanderbricks.payments;
Para más información sobre cómo combinar un relleno puntual con CDC en curso en el mismo objetivo, véase Relleno de datos históricos con oleoductos.