Pular para o conteúdo principal

CRIAR FLUXO (pipeline)

Use a instrução CREATE FLOW para criar fluxos ou preenchimentos retroativos para tabelas em um pipeline.

nota

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.

  • AUTO CDC EM

    Uma instrução AUTO CDC ... INTO que define o fluxo, com um create_auto_cdc_flow_spec. Você deve incluir uma declaração AUTO CDC ... INTO ou uma declaração INSERT INTO . Use AUTO CDC ... INTO quando 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 ... INTO que 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ória FROM SNAPSHOT (snapshot_query) que lê os dados de snapshot e uma cláusula opcional WITH 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 usando KEYS para a identidade da linha e STORED AS para 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 (...). Se WITH VERSION (...) não for especificado, current_snapshot_version() não poderá ser chamado dentro de FROM SNAPSHOT (...).

      Quando WITH VERSION (...) é omitido, o mecanismo lê a origem diretamente por meio de FROM 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á com AUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. Para processar Snapshot em várias atualizações, 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.

      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 with AUTO_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 WHERE ou SEQUENCE BY. A ordenação entre Snapshot é expressa por meio de WITH 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 ONCE nã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ção skipChangeCommits para lidar com erros.

    INSERT INTO é mutuamente exclusivo com AUTO CDC ... INTO. Use AUTO CDC ... INTO quando os dados de origem incluem a funcionalidade de captura de dados de alterações (CDC) (CDC). Use INSERT INTO quando 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

info

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 a condition sã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 a condition nã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 WHERE uses lotes semantics, so the source query doesn't need to be a transmissão query. BY NAME is required. REPLACE WHERE não pode ser combinado com ONCE, REPLACE USING ou AUTO 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 ONCE altera o fluxo de duas maneiras:

    • A fonte query ou create_auto_cdc_flow_spec nã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.

    ONCE não pode ser usado com REPLACE USING, que requer uma fonte de transmissão.

Exemplos​

SQL
-- 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.