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.
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.
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
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
)
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 |
|---|---|---|
|
| 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 |
|
| 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 |
|
| 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:
Argumentos para funções |
|
| Se os registros devem ser armazenados como SCD tipo 1 ou SCD tipo 2. Defina como |
|
| Um subconjunto de colunas de saída a serem rastreadas para histórico na tabela de destino. Use
Argumentos para funções |
|
| Indica se o fluxo de snapshot deve ser executado apenas uma vez. Quando |
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 CDCflow 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 CDCkey names. MultipleAUTO CDCflows 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 definaignore_null_updates,ignore_null_updates_column_listouignore_null_updates_except_column_listnos fluxoscreate_auto_cdc_flow(). - Se você adicionar um fluxo
AUTO CDCa um destinoAUTO CDC FROM SNAPSHOTexistente, o destino não deve conter colunas de usuário cujos nomes entrem em conflito com as colunas de sistema reservadasAUTO 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:
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 é
Noneou 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:
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:
- execute a função
next_snapshot_and_versionpara carregar o próximo Snapshot DataFrame e a versão correspondente do Snapshot. - 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.
- Detecta as alterações no novo Snapshot e as aplica incrementalmente à tabela de destino.
- Retorna ao passo #1 para carregar o próximo Snapshot e sua versão.