Pular para o conteúdo principal

criar_auto_cdc_a_partir_do_fluxo_de_instantâneo

A função create_auto_cdc_from_snapshot_flow cria um fluxo que usa a funcionalidade de captura de dados de alterações (CDC) dos LakeFlow Pipelines para processar dados de origem de Snapshots de banco de dados. Consulte Como o AUTO CDC FROM SNAPSHOT funciona.

nota

Esta função substitui a função anterior apply_changes_from_snapshot(). As duas funções têm a mesma assinatura. A Databricks recomenda atualizar para usar o novo nome.

importante

Você deve ter uma tabela de transmissão de destino para esta operação. Para criar a tabela de destino necessária, você pode usar a função create_streaming_table().

Sintaxe​

Python
from pyspark import pipelines as dp

dp.create_auto_cdc_from_snapshot_flow(
target = "<target-table>",
source = Any,
keys = ["key1", "key2", "keyN"],
stored_as_scd_type = "1",
track_history_column_list = None,
track_history_except_column_list = None,
once = False
)
nota

Para o processamento AUTO CDC FROM SNAPSHOT , o comportamento default é inserir uma nova linha quando um registro correspondente com a(s) mesma(s) key(s) não existe(m) no destino. Se existir um registro correspondente, ele será atualizado somente se algum dos valores na linha tiver sido alterado. Linhas com chave presente no destino, mas não mais presente na origem, são excluídas.

Para saber mais sobre o processamento de CDC com snapshots, consulte As APIs do AUTO CDC: Simplifique a captura de dados de alterações (CDC) com pipelines. Para ver exemplos de uso da função create_auto_cdc_from_snapshot_flow(), consulte os exemplos de ingestão periódica de snapshot, ingestão histórica de snapshot e backfill do SCD Tipo 1.

Parâmetros​

Parâmetro

Tipo

Descrição

target

str

Obrigatório. O nome da tabela a ser atualizada. Você pode usar a função create_streaming_table() para criar a tabela de destino antes de executar a função create_auto_cdc_from_snapshot_flow() .

source

str ou lambda function

Obrigatório. O nome de uma tabela ou view para Snapshot periodicamente ou uma função lambda Python que retorna o Snapshot DataFrame a ser processado e a versão do Snapshot. Veja Implementar o argumento source.

keys

list

Obrigatório. A coluna ou combinação de colunas que identifica exclusivamente uma linha nos dados de origem. Isso é usado para identificar quais eventos do CDC se aplicam a registros específicos na tabela de destino. Você pode especificar:

  • Uma lista de strings: ["userId", "orderId"]

  • Uma lista de funções Spark SQL col() : [col("userId"), col("orderId"].

Argumentos para funções col() não podem incluir qualificadores. Por exemplo, você pode usar col(userId), mas não pode usar col(source.userId).

stored_as_scd_type

str ou int

Se os registros devem ser armazenados como SCD tipo 1 ou SCD tipo 2. Defina como 1 para SCD tipo 1 ou 2 para SCD tipo 2. O default é SCD tipo 1.

track_history_column_list ou track_history_except_column_list

list

Um subconjunto de colunas de saída a serem rastreadas para histórico na tabela de destino. Use track_history_column_list para especificar a lista completa de colunas a serem rastreadas. Use track_history_except_column_list para especificar as colunas a serem excluídas do acompanhamento. Você pode declarar qualquer valor como uma lista de strings ou como funções Spark SQL col() :

  • track_history_column_list = ["userId", "name", "city"]
  • track_history_column_list = [col("userId"), col("name"), col("city")]
  • track_history_except_column_list = ["operation", "sequenceNum"]
  • track_history_except_column_list = [col("operation"), col("sequenceNum")

Argumentos para funções col() não podem incluir qualificadores. Por exemplo, você pode usar col(userId), mas não pode usar col(source.userId). O default é incluir todas as colunas na tabela de destino quando nenhum argumento track_history_column_list ou track_history_except_column_list é passado para a função.

once

bool

Indica se o fluxo de snapshot deve ser executado apenas uma vez. Quando True, o fluxo é ignorado durante atualizações incrementais subsequentes após o commit bem-sucedido. Uma refresh completa do destino reexecuta o fluxo. Use esta opção para um carregamento único de snapshot ou preenchimento retroativo. O default é False.

Parâmetro

Tipo

Descrição

target

str

Obrigatório. O nome da tabela a ser atualizada. Você pode usar a função create_streaming_table() para criar a tabela de destino antes de executar a função create_auto_cdc_from_snapshot_flow() .

source

str ou lambda function

Obrigatório. O nome de uma tabela ou view para Snapshot periodicamente ou uma função lambda Python que retorna o Snapshot DataFrame a ser processado e a versão do Snapshot. Veja Implementar o argumento source.

keys

list

Obrigatório. A coluna ou combinação de colunas que identifica exclusivamente uma linha nos dados de origem. Isso é usado para identificar quais eventos do CDC se aplicam a registros específicos na tabela de destino. Você pode especificar:

  • Uma lista de strings: ["userId", "orderId"]

  • Uma lista de funções Spark SQL col() : [col("userId"), col("orderId"].

Argumentos para funções col() não podem incluir qualificadores. Por exemplo, você pode usar col(userId), mas não pode usar col(source.userId).

stored_as_scd_type

str ou int

Se os registros devem ser armazenados como SCD tipo 1 ou SCD tipo 2. Defina como 1 para SCD tipo 1 ou 2 para SCD tipo 2. O default é SCD tipo 1.

track_history_column_list ou track_history_except_column_list

list

Um subconjunto de colunas de saída a serem rastreadas para histórico na tabela de destino. Use track_history_column_list para especificar a lista completa de colunas a serem rastreadas. Use track_history_except_column_list para especificar as colunas a serem excluídas do acompanhamento. Você pode declarar qualquer valor como uma lista de strings ou como funções Spark SQL col() :

  • track_history_column_list = ["userId", "name", "city"]
  • track_history_column_list = [col("userId"), col("name"), col("city")]
  • track_history_except_column_list = ["operation", "sequenceNum"]
  • track_history_except_column_list = [col("operation"), col("sequenceNum")

Argumentos para funções col() não podem incluir qualificadores. Por exemplo, você pode usar col(userId), mas não pode usar col(source.userId). O default é incluir todas as colunas na tabela de destino quando nenhum argumento track_history_column_list ou track_history_except_column_list é passado para a função.

once

bool

Indica se o fluxo de snapshot deve ser executado apenas uma vez. Quando True, o fluxo é ignorado durante atualizações incrementais subsequentes após o commit bem-sucedido. Uma refresh completa do destino reexecuta o fluxo. Use esta opção para um carregamento único de snapshot ou preenchimento retroativo. O default é False.

Notas​

Para destinos SCD tipo 1 em pipelines com trigger, um fluxo create_auto_cdc_from_snapshot_flow() pode compartilhar um destino com um ou mais fluxos AUTO CDC definidos em Python ou SQL. Aplicam-se os seguintes requisitos:

  • Give each AUTO CDC flow a unique name.
  • Use the same number of keys in the same order for all flows. Snapshot-flow key names are compared case-insensitively with AUTO CDC key names. Multiple AUTO CDC flows must use identical key names and casing.
  • Use exatamente o mesmo tipo de dados para a versão de snapshot e a coluna de sequenciamento de cada fluxo de AUTO CDC.
  • Não defina expectativas no fluxo AUTO CDC FROM SNAPSHOT.
  • Não use IGNORE NULL UPDATES. Em Python, não defina ignore_null_updates, ignore_null_updates_column_list ou ignore_null_updates_except_column_list nos fluxos create_auto_cdc_flow().
  • Se você adicionar um fluxo AUTO CDC a um destino AUTO CDC FROM SNAPSHOT existente, o destino não deve conter colunas de usuário cujos nomes entrem em conflito com as colunas de sistema reservadas AUTO CDC.
  • Não use esse padrão para o SCD Tipo 2 ou destinos bitemporais.

Para ver um exemplo, consulte Add a backfill to an AUTO CDC SCD Type 1 table.

Implementar o argumento source​

A função create_auto_cdc_from_snapshot_flow() inclui o argumento source . Para processar o Snapshot histórico, espera-se que o argumento source seja uma função lambda Python que retorna dois valores para a função create_auto_cdc_from_snapshot_flow() : um DataFrame Python contendo os dados do Snapshot a serem processados e uma versão do Snapshot.

A primeira invocação da função deve retornar um DataFrame de snapshot e a versão. Se a primeira invocação retornar None, a atualização do pipeline falhará. Após o processamento de pelo menos um snapshot, retorne None para indicar que não há snapshots adicionais disponíveis.

A seguir está a assinatura da função lambda:

Python
lambda Any => Optional[(DataFrame, Any)]
  • O argumento para a função lambda é a versão mais recentemente processada do Snapshot.
  • O valor de retorno da função lambda é None ou uma tupla de dois valores: O primeiro valor da tupla é um DataFrame contendo o Snapshot a ser processado. O segundo valor da tupla é a versão do Snapshot que representa a ordem lógica do Snapshot.

Um exemplo que implementa e chama a função lambda:

Python
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Tuple[DataFrame, Optional[int]]:
if latest_snapshot_version is None:
return (spark.read.load("filename.csv"), 1)
else:
return None

create_auto_cdc_from_snapshot_flow(
# ...
source = next_snapshot_and_version,
# ...
)

O runtime do Lakeflow pipelines executa os seguintes passos cada vez que o pipeline que contém a função create_auto_cdc_from_snapshot_flow() é acionado:

  1. execute a função next_snapshot_and_version para carregar o próximo Snapshot DataFrame e a versão correspondente do Snapshot.
  2. Se a primeira invocação não retornar nenhum DataFrame, a atualização falhará. Se uma invocação posterior não retornar nenhum DataFrame, a execução será encerrada e a atualização do pipeline será marcada como concluída.
  3. Detecta as alterações no novo Snapshot e as aplica incrementalmente à tabela de destino.
  4. Retorna ao passo #1 para carregar o próximo Snapshot e sua versão.