CRIAR FLUXO (pipelines)

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

Sintaxe

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

    O nome do fluxo a ser criado.

  • COMENTAR

    Uma descrição opcional para o fluxo.

  • AUTO CDC PARA

    Uma AUTO CDC ... INTO declaração que define o fluxo, com um create_auto_cdc_flow_spec. Você deve incluir uma AUTO CDC ... INTO declaração ou uma INSERT INTO declaração. Use AUTO CDC ... INTO quando a consulta de origem usa a semântica de dados de alteração.

    Para obter mais informações, consulte AUTO CDC INTO (pipelines).

  • 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 ONCE opçã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ção skipChangeCommits para manipular erros.

    INSERT INTO é mutuamente exclusivo com AUTO CDC ... INTO. Use AUTO CDC ... INTO quando os dados de origem incluírem a funcionalidade de captura de dados de alteração (CDC). Use INSERT INTO quando 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 USING fluxo, que substitui todas as linhas da tabela de destino que correspondam às colunas-chave especificadas e deixa todas as outras linhas intocadas. Use REPLACE USING quando a 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 vença, mesmo quando as atualizações chegam fora de ordem.

    Especifique pelo menos uma coluna de chave e exatamente uma SEQUENCE BY coluna. A consulta deve ser uma consulta de streaming e BY NAME é obrigatória. REPLACE USING não pode ser combinado com ONCE ou com AUTO CDC ... INTO.

    Para mais informações, veja Substituição parcial de snapshot por SUBSTITUIR os fluxos USANDO.

  • UMA VEZ

    Opcionalmente, defina o fluxo como um fluxo único, como um preenchimento adicional. O uso ONCE altera o fluxo de duas maneiras:

    • A fonte query ou create_auto_cdc_flow_spec não é uma tabela de streaming.
    • O fluxo é executado uma vez por padrão. Se o pipeline receber uma atualização completa, o fluxo ONCE será executado novamente para recriar os dados.

    ONCE Não pode ser usado com REPLACE USING, que requer uma fonte 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);