CREATE FLOW (pipelines)

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

Note

CREATE FLOWDirecionar uma tabela de streaming suporta tanto AUTO CDC ... INTO os fluxos U.REPLACE WHERE Uma tabela gerenciada criada com CREATE TABLE ... FLOW não suporta captura de dados de alteração: um AUTO CDC ... INTO fluxo contra uma tabela gerenciada falha com MANAGED_TABLE_DOES_NOT_SUPPORT_CDC. CREATE FLOW Into a Streaming Table não está sujeito a essa limitação. Veja CREATE TABLE ... FLOW (tubulações).

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.

  • COMENTÁRIO

    Uma descrição opcional para o fluxo.

  • AUTO-CDC EM

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

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

  • AUTO CDC ... DO SNAPSHOT

    Uma AUTO CDC ... INTO instrução que deriva as alterações comparando instantâneos em vez de ler um feed de mudanças. Use este formulário quando a captura de dados de alteração não estiver habilitada na fonte e apenas capturas completas estiverem disponíveis. A fonte é especificada em duas partes: uma cláusula obrigatória FROM SNAPSHOT (snapshot_query) que lê os dados do snapshot e uma cláusula opcional WITH VERSION (version_query) que seleciona a próxima versão do snapshot a ser processada. Veja como o AUTO CDC FROM SNAPSHOT funciona.

    • DE SNAPSHOT (snapshot_query)

      Required. Uma consulta que lê os dados 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 deleções, e os funde no destino usando KEYS a identidade da linha e STORED AS para determinar como as mudanças são armazenadas.

      Chame current_snapshot_version() dentro desta consulta para referenciar a versão selecionada por WITH VERSION (...). Se WITH VERSION (...) não for especificado, current_snapshot_version() não pode ser chamada dentro FROM SNAPSHOT (...)de .

      Quando WITH VERSION (...) é omitido, o motor lê a fonte diretamente através FROM SNAPSHOT (...)de , e a consulta snapshot roda apenas durante a carga inicial, enquanto o destino não possui 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 com AUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. Para processar snapshots em múltiplas atualizações, use WITH VERSION (...).

    • COM VERSÃO (version_query)

      Optional. Uma consulta que seleciona a próxima versão snapshot a ser processada. Ele deve retornar exatamente uma coluna de tipo ordenável e 0 ou 1 linha. Quando retorna 1 linha, o valor deve ser não nulo. A coluna pode ser um valor escalar, como um BIGINT, ou a STRUCT cujos campos são todos ordenáveis. Uma consulta de versão que retorna mais de uma coluna, mais de uma linha ou um valor nulo falha no fluxo com INVALID_AUTO_CDC_FROM_SNAPSHOT_VERSION_QUERY.

      Durante uma atualização do pipeline, o motor repete os seguintes passos: ele avalia a consulta de versão; se a consulta devolver 0 linhas, para de processar esse fluxo para a atualização atual; se a consulta retornar 1 linha, o motor expõe esse valor por meio 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 então reavalia a consulta de versão para selecionar a próxima versão. Uma única atualização de pipeline processa as versões em ordem até que a consulta de versão não retorne linhas.

      Toda versão retornada após um commit bem-sucedido deve ser maior que a versão previamente confirmada; uma versão não crescente falha na 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 mudança no tipo de dado falha na atualização com AUTO_CDC_FROM_SNAPSHOT_VERSION_SCHEMA_CHANGED. Uma atualização completa elimina o estado da versão persistente.

    • CHAVES

      Required. As colunas principais são usadas para identificar linhas entre instantâneos para detecção de mudanças.

    • STORED AS { SCD TYPE 1 | SCD TYPE 2 }

      Optional. Especifica como as mudanças são armazenadas na tabela de destino. O padrão é SCD TYPE 1.

    • TRACK HISTORY ON { col_list | * EXCEPT (col_list) }

      Optional. Vale apenas com SCD TYPE 2. Especifica quais colunas acionam uma nova linha de histórico quando elas mudam. Forneça uma lista explícita de colunas ou * EXCEPT (col_list) rastreie todas as colunas, exceto as listadas.

    Snapshot CDC não suporta WHERE nem SEQUENCE BY. A ordenação cruzada de snapshots é expressa através de WITH VERSION (...).

  • target_table

    A tabela a ser atualizada. Essa deve ser uma tabela streaming.

  • INSERT EM

    Define uma consulta de tabela inserida na tabela de destino. Se a opção ONCE não for fornecida, a consulta deverá ser uma consulta de streaming . Use a palavra-chave STREAM para usar a semântica de streaming 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ções de alterações, 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 incluírem a funcionalidade de CDC (captura de dados de alteração). Use INSERT INTO quando a origem não o fizer.

    Para obter mais informações sobre dados de fluxo, consulte Transformar dados com pipelines.

  • SUBSTITUA USANDO ( column_name [, ...] ) SEQUÊNCIA POR sequence_column

    Importante

    Esse recurso está em 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 correspondem às colunas-chave especificadas e deixa todas as outras linhas intocadas. 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 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 snapshots por SUBSTITUIR fluxos USANDO.

  • Condição SUBSTITUIR WHERE

    Define o fluxo como um REPLACE WHERE fluxo, que recalcula e sobrescreve um subconjunto alvo da tabela alvo. A cada atualização, todas as linhas da correspondência condition da tabela alvo são deletadas, a consulta de origem é recalculada para o mesmo intervalo de predicados e os resultados são inseridos. As fileiras que não combinam condition ficam intocadas. Você não precisa adicionar o predicado à consulta de origem; o motor de pipeline o aplica automaticamente ao ler da fonte.

    REPLACE WHERE Usa semântica em lote, então a consulta de origem não precisa ser uma consulta de streaming. BY NAME é obrigatório. REPLACE WHERE Não pode ser combinado com ONCE, REPLACE USING, ou AUTO CDC ... INTO.

    Para mais informações, veja 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 backfill. Usar ONCE altera o fluxo de duas maneiras:

    • A origem query ou create_auto_cdc_flow_spec não é uma tabela de streaming.
    • O fluxo é executado uma vez por padrão. Se a linha de processamento for atualizada com um recarregamento completo, então 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.

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 backfill único com CDC em andamento no mesmo alvo, veja Backfill de dados históricos com pipelines.