Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Use a CREATE FLOW instrução para criar fluxos ou backfills para tabelas em um pipeline.
Note
CREATE FLOW Direcionar uma tabela de streaming suporta ambos AUTO CDC ... INTO os fluxos AND REPLACE WHERE . Uma tabela gerida criada com CREATE TABLE ... FLOW não suporta captura de dados de alteração: um AUTO CDC ... INTO fluxo contra uma tabela gerida falha com MANAGED_TABLE_DOES_NOT_SUPPORT_CDC.
CREATE FLOW Entrar numa mesa de streaming não está sujeito a essa limitação. Veja CREATE TABLE ... FLOW (oleodutos).
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
flow_name
O nome do fluxo a ser criado.
COMENTAR
Uma descrição opcional para o fluxo.
-
Uma
AUTO CDC ... INTOdeclaração que define o fluxo, com umcreate_auto_cdc_flow_spec. Você deve incluir umaAUTO CDC ... INTOdeclaração ou umaINSERT INTOdeclaração. UseAUTO CDC ... INTOquando a consulta de origem usa a semântica de dados de alteração.Para obter mais informações, consulte AUTO CDC INTO (pipelines).
AUTO CDC ... DO INSTANTÂNEO
Uma
AUTO CDC ... INTOafirmação que deriva alterações comparando instantâneos em vez de ler um feed de alterações. Use este formulário quando a captura de dados de alteração não estiver ativada na fonte e apenas estiverem disponíveis instantâneos completos. A fonte é especificada em duas partes: uma cláusula obrigatóriaFROM SNAPSHOT (snapshot_query)que lê os dados do snapshot e uma cláusula opcionalWITH VERSION (version_query)que seleciona a próxima versão do snapshot a processar. Veja como funciona o AUTO CDC FROM SNAPSHOT.DE SNAPSHOT (snapshot_query)
Required. Uma consulta que lê os dados do snapshot para a versão selecionada por
WITH VERSION (...). O motor diferencia o resultado do snapshot previamente comprometido para derivar inserções, atualizações e eliminações, e funde-os no destino usandoKEYSa identidade da linha eSTORED ASpara determinar como as alterações são armazenadas.Chame current_snapshot_version() dentro desta consulta para referenciar a versão selecionada por
WITH VERSION (...). SeWITH VERSION (...)não for especificado,current_snapshot_version()não é chamável dentroFROM SNAPSHOT (...)de .Quando
WITH VERSION (...)é omitido, o motor lê a fonte diretamente atravésFROM SNAPSHOT (...)de , e a consulta snapshot executa-se apenas durante a carga inicial, enquanto o destino não tem dados comprometidos nem estado de snapshot comprometido. Em qualquer atualização posterior, quando o destino já contém dados ou tem estado de snapshot comprometido, o fluxo falha comAUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. Para processar instantâneos ao longo de múltiplas atualizações, useWITH VERSION (...).COM VERSÃO (version_query)
Optional. Uma consulta que seleciona a próxima versão snapshot a processar. Deve devolver exatamente uma coluna de tipo ordenável e 0 ou 1 linha. Quando devolve 1 linha, o valor deve ser não nulo. A coluna pode ser um valor escalar, como um
BIGINT, ou aSTRUCTcujos campos são todos ordenáveis. Uma consulta de versão que devolve mais do que uma coluna, mais de uma linha ou um valor nulo falha o fluxo comINVALID_AUTO_CDC_FROM_SNAPSHOT_VERSION_QUERY.Durante uma atualização do pipeline, o motor repete os seguintes passos: avalia a consulta de versão; se a consulta devolver 0 linhas, para de processar este fluxo para a atualização atual; se a consulta devolver 1 linha, o motor expõe esse valor através de current_snapshot_version(), avalia a consulta snapshot, compromete o snapshot resultante e expõe a versão comprometida através de last_snapshot_version(). O motor reavalia então a consulta de versão para selecionar a versão seguinte. Uma única atualização de pipeline processa as versões por ordem até que a consulta de versão não devolva linhas.
Cada versão devolvida após um commit bem-sucedido deve ser maior do que a versão previamente confirmada; uma versão não crescente falha a atualização com
APPLY_CHANGES_FROM_SNAPSHOT_ERROR.OUT_OF_ORDER_SNAPSHOT_VERSION. O tipo de dado do valor da versão deve permanecer inalterado entre os commits snapshot; uma alteração do tipo de dado falha a atualização comAUTO_CDC_FROM_SNAPSHOT_VERSION_SCHEMA_CHANGED. Uma atualização completa limpa o estado da versão persistente.CHAVES
Required. As principais colunas-chave são usadas para identificar linhas entre instantâneos para deteção de alterações.
STORED AS { SCD TYPE 1 | SCD TYPE 2 }Optional. Especifica como as alterações são armazenadas na tabela alvo. A predefinição é
SCD TYPE 1.TRACK HISTORY ON { col_list | * EXCEPT (col_list) }Optional. Aplica-se apenas com
SCD TYPE 2. Especifica quais as colunas que acionam uma nova linha de histórico quando mudam. Forneça uma lista explícita de colunas ou* EXCEPT (col_list)para acompanhar todas as colunas exceto as listadas.
Snapshot CDC não suporta
WHEREnemSEQUENCE BY. A ordenação cruzada de instantâneos é expressa através deWITH VERSION (...).target_table
A tabela a ser atualizada. Esta deve ser uma tabela de Streaming.
INSERT EM
Define uma consulta de tabela que é inserida na tabela de destino. Se a
ONCEopção não for fornecida, a consulta deverá ser uma consulta de streaming . Utilize a palavra-chave STREAM para aplicar a semântica de transmissão e ler a partir da fonte. Se a leitura encontrar uma alteração ou exclusão em um registro existente, um erro será gerado. É mais seguro ler a partir de fontes estáticas ou apenas de anexação. Para ingerir dados que têm confirmações de alteração, podes usar Python e a opçãoskipChangeCommitspara manipular erros.INSERT INTOé mutuamente exclusivo comAUTO CDC ... INTO. UseAUTO CDC ... INTOquando os dados de origem incluírem a funcionalidade de captura de dados de alteração (CDC). UseINSERT INTOquando a fonte não o fizer.Para obter mais informações sobre a transmissão de dados, consulte Transformar dados com pipelines.
SUBSTITUIR USANDO ( column_name [, ...] ) SEQUÊNCIA POR sequence_column
Importante
Este recurso está em versão Beta. Requer Databricks Runtime 18.2 e superiores.
Define o fluxo como um
REPLACE USINGfluxo, que substitui todas as linhas da tabela de destino que correspondam às colunas-chave especificadas e deixa todas as outras linhas intocadas. UseREPLACE USINGquando a sua fonte for uma série de snapshots parciais indexados por coluna.SEQUENCE BYordena as atualizações para que a sequência mais alta para uma chave vença, mesmo quando as atualizações chegam fora de ordem.Especifique pelo menos uma coluna de chave e exatamente uma
SEQUENCE BYcoluna. A consulta deve ser uma consulta de streaming eBY NAMEé obrigatória.REPLACE USINGnão pode ser combinado comONCEou comAUTO CDC ... INTO.Para mais informações, veja Substituição parcial de snapshot por SUBSTITUIR os fluxos USANDO.
Condição REPLACE WHERE
Define o fluxo como um
REPLACE WHEREfluxo, que recalcula e sobrescreve um subconjunto alvo da tabela alvo. Em cada atualização, todas as linhas da correspondênciaconditionda tabela alvo são eliminadas, a consulta de origem é recalculada para esse mesmo intervalo de predicados e os resultados são inseridos. As filas que não coincidemconditionficam intocadas. Não precisas de adicionar o predicado à consulta de origem; o motor de pipeline aplica-o automaticamente ao ler da fonte.REPLACE WHEREUsa semântica em lote, por isso a consulta de origem não precisa de ser uma consulta de streaming.BY NAMEé obrigatório.REPLACE WHEREnão pode ser combinado comONCE,REPLACE USING, ouAUTO CDC ... INTO.Para mais informações, consulte processamento em lote com fluxos REPLACE WHERE e fluxos REPLACE WHERE para tabelas de streaming independentes.
UMA VEZ
Opcionalmente, defina o fluxo como um fluxo único, como um preenchimento adicional. O uso
ONCEaltera o fluxo de duas maneiras:- A fonte
queryoucreate_auto_cdc_flow_specnão é uma tabela de streaming. - O fluxo é executado uma vez por padrão. Se o pipeline receber uma atualização completa, o fluxo
ONCEserá executado novamente para recriar os dados.
ONCENão pode ser usado comREPLACE USING, que requer uma fonte de streaming.- A fonte
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 saber mais sobre a combinação de um backfill único com o CDC em curso no mesmo objetivo, veja Backfill de dados históricos com pipelines.