CRIAR FLUXO (pipeline)
Use a instrução CREATE FLOW para criar fluxos ou preenchimentos retroativos para tabelas em um pipeline.
CREATE FLOW o direcionamento para uma tabela de transmissão é compatível com fluxos AUTO CDC ... INTO e REPLACE WHERE. Uma tabela gerenciada criada com CREATE TABLE ... FLOW não é compatível com a captura de dados de alterações (CDC): um fluxo AUTO CDC ... INTO em uma tabela gerenciada falha com MANAGED_TABLE_DOES_NOT_SUPPORT_CDC. CREATE FLOW em uma tabela de transmissão não está sujeito a essa limitação. Consulte CREATE TABLE ... FLOW (pipelines).
Sintaxe
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
-
nome_do_fluxo
O nome do fluxo a ser criado.
-
comentário
Uma descrição opcional para o fluxo.
-
Uma instrução
AUTO CDC ... INTOque define o fluxo, com umcreate_auto_cdc_flow_spec. Você deve incluir uma declaraçãoAUTO CDC ... INTOou uma declaraçãoINSERT INTO. UseAUTO CDC ... INTOquando a consulta de origem usar semântica de dados alterados.Para obter mais informações, consulte AUTO CDC INTO (pipeline).
-
AUTO CDC ... FROM SNAPSHOT
Uma instrução
AUTO CDC ... INTOque deriva alterações comparando snapshots em vez de ler um feed de alterações. Use este formulário quando a captura de dados de alterações (CDC) não estiver habilitada na origem e apenas Snapshots completos estiverem disponíveis. A origem é especificada em duas partes: uma cláusula obrigatóriaFROM SNAPSHOT (snapshot_query)que lê os dados de snapshot e uma cláusula opcionalWITH VERSION (version_query)que seleciona a próxima versão do snapshot a ser processada. Consulte Como funciona o AUTO CDC FROM SNAPSHOT.-
FROM SNAPSHOT (snapshot_query)
Obrigatório. Uma query que lê os dados do snapshot para a versão selecionada por
WITH VERSION (...). O engine compara o resultado com o snapshot previamente confirmado para derivar inserções, atualizações e exclusões, e as faz merge no destino usandoKEYSpara a identidade da linha eSTORED ASpara determinar como as alterações são armazenadas.Chame current_snapshot_version() dentro desta query para fazer referência à versão selecionada por
WITH VERSION (...). SeWITH VERSION (...)não for especificado,current_snapshot_version()não poderá ser chamado dentro deFROM SNAPSHOT (...).Quando
WITH VERSION (...)é omitido, o mecanismo lê a origem diretamente por meio deFROM SNAPSHOT (...), e a query de snapshot é executada apenas durante a carga inicial, enquanto a tabela de destino não tiver dados consolidados nem estado de snapshot consolidado. Em qualquer atualização posterior, quando a tabela de destino já contiver dados ou tiver um estado de snapshot consolidado, o fluxo falhará comAUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. Para processar Snapshot em várias atualizações, 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.Durante uma atualização de pipeline, o engine repete os seguintes passos: ele avalia a query de versão; se a query retornar 0 linhas, ele para de processar este fluxo para a atualização atual; se a query retornar 1 linha, o engine expõe esse valor por meio de current_snapshot_version(), avalia a query de Snapshot, faz o commit do Snapshot resultante e expõe a versão com commit por meio de last_snapshot_version(). O engine reavalia então a query de versão para selecionar a próxima versão. Uma única atualização de pipeline processa versões em ordem até que a query de versão não retorne nenhuma linha.
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. -
key
Obrigatório. As colunas de key primária usadas para identificar linhas em snapshots para detecção de alterações.
-
STORED AS { SCD TYPE 1 | SCD TYPE 2 }Opcional. Especifica como as alterações são armazenadas na tabela de destino. O default é
SCD TYPE 1. -
TRACK HISTORY ON { col_list | * EXCEPT (col_list) }Opcional. Aplica-se apenas com
SCD TYPE 2. Especifica quais colunas Trigger uma nova linha de história quando são alteradas. Forneça uma lista de colunas explícita ou* EXCEPT (col_list)para rastrear todas as colunas, exceto as listadas.
Snapshot CDC não oferece suporte a
WHEREouSEQUENCE BY. A ordenação entre Snapshot é expressa por meio deWITH VERSION (...). -
-
tabela_alvo
A tabela a ser atualizada. Esta deve ser uma tabela de transmissão.
-
INSERT INTO
Define uma consulta de tabela que é inserida na tabela de destino. Se a opção
ONCEnão for fornecida, a consulta deverá ser uma consulta de transmissão . Use a palavra-chave transmissão para usar a semântica de transmissão para ler a fonte. Se a leitura encontrar uma alteração ou exclusão em um registro existente, um erro será gerado. É mais seguro ler de fontes estáticas ou somente de acréscimos. Para ingerir dados que tenham confirmação de alteração, você pode usar Python e a opçãoskipChangeCommitspara lidar com erros.INSERT INTOé mutuamente exclusivo comAUTO CDC ... INTO. UseAUTO CDC ... INTOquando os dados de origem incluem a funcionalidade de captura de dados de alterações (CDC) (CDC). UseINSERT INTOquando a fonte não o fizer.Para mais informações sobre transmissão de dados, veja transformação de dados com pipeline.
-
REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column
Beta
Esse recurso está em Beta. Requer Databricks Runtime 18.2 ou acima.
Define o fluxo como um fluxo REPLACE USING, que substitui todas as linhas na tabela de destino que correspondem às colunas de key especificadas e deixa todas as outras linhas inalteradas. Use REPLACE USING quando sua fonte for uma série de Snapshots parciais indexados por coluna. SEQUENCE BY ordena as atualizações para que a sequência mais alta para uma chave prevaleça, mesmo quando as atualizações chegam fora de ordem.
Especifique pelo menos uma coluna key e exatamente uma coluna SEQUENCE BY. A query deve ser uma query de transmissão, e BY NAME é obrigatório. REPLACE USING não pode ser combinado com ONCE ou com AUTO CDC ... INTO.
Para obter mais informações, consulte Substituição de snapshot parcial com fluxos REPLACE USING.
-
Condição REPLACE WHERE
Define o fluxo como um fluxo
REPLACE WHERE, que recomputa e sobrescreve um subconjunto direcionado da tabela de destino. Em cada atualização, todas as linhas na tabela de destino que correspondem aconditionsão excluídas, a query de origem é recalculada para esse mesmo intervalo de predicados e os resultados são inseridos. As linhas que não correspondem aconditionnão são modificadas. Você não precisa adicionar o predicado à query de origem; o mecanismo do pipeline o aplica automaticamente ao ler a partir da origem.REPLACE WHEREuses lotes semantics, so the source query doesn't need to be a transmissão query.BY NAMEis required.REPLACE WHEREnão pode ser combinado comONCE,REPLACE USINGouAUTO CDC ... INTO.For more informação, see lotes processing with REPLACE WHERE flows and REPLACE WHERE flows for standalone transmissão tables.
-
UMA VEZ
Opcionalmente, defina o fluxo como um fluxo único, como um aterro. Usar
ONCEaltera o fluxo de duas maneiras:- A fonte
queryoucreate_auto_cdc_flow_specnão é uma tabela de transmissão. - O fluxo é executado uma vez por padrão. Se o pipeline for atualizado com um refresh completo, então o fluxo
ONCEé executado novamente para recriar os dados.
ONCEnão pode ser usado comREPLACE USING, que requer uma fonte de transmissão. - A fonte
Exemplos
-- 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 saber mais sobre como combinar um preenchimento único com o CDC contínuo no mesmo destino, consulte Preenchendo data histórica com pipelines.