Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
A create_auto_cdc_from_snapshot_flow função cria um fluxo que usa a funcionalidade CDC (captura de dados) de pipelines do Lakeflow para processar dados de origem de instantâneos de banco de dados. Veja como o AUTO CDC FROM SNAPSHOT funciona.
Observação
Essa função substitui a função apply_changes_from_snapshot()anterior. As duas funções têm a mesma assinatura. O Databricks recomenda a atualização para usar o novo nome.
Importante
Você deve possuir uma tabela de destino de streaming 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
)
Observação
Para AUTO CDC FROM SNAPSHOT processamento, o comportamento padrão é inserir uma nova linha quando um registro correspondente com a(s) mesma(s) chave(s) não existe no destino. Se um registro correspondente existir, ele será atualizado somente se algum dos valores na linha tiver sido alterado. As linhas com chaves presentes no destino, mas não mais presentes na origem, são excluídas.
Para saber mais sobre o processamento CDC com instantâneos, consulte as APIs AUTO CDC: Simplifique a captura de dados de mudanças com pipelines. Para exemplos de uso da create_auto_cdc_from_snapshot_flow() função, veja os exemplos de ingestão periódica de snapshots, ingestão histórica de snapshots e backfill SCD Tipo 1 .
Parâmetros
| Parâmetro | Tipo | Description |
|---|---|---|
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 create_auto_cdc_from_snapshot_flow() função. |
source |
str ou lambda function |
Obrigatório O nome de uma tabela ou visualização para capturar instantâneos periodicamente ou uma função lambda do Python que retorna o DataFrame do instantâneo a ser processado e a versão do instantâneo. Consulte Implementar o source argumento. |
keys |
list |
Obrigatório A coluna ou combinação de colunas que identificam exclusivamente uma linha nos dados de origem. Isso é usado para identificar quais eventos CDC se aplicam a registros específicos na tabela de destino. Você pode especificar:
|
stored_as_scd_type |
str ou int |
Se deseja armazenar registros como SCD tipo 1 ou SCD tipo 2. Defina para SCD 1 tipo 1 ou 2 SCD tipo 2. O padrão é SCD tipo 1. |
track_history_column_list ou track_history_except_column_list |
list |
Um subconjunto das colunas de saída a ser registrado no 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 o valor como uma lista de cadeias de caracteres ou como funções SQL col() do Spark:
Argumentos para col() funções não podem incluir qualificadores. Por exemplo, você pode usar col(userId), mas não pode usar col(source.userId). O padrão é 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 |
Se deve rodar o fluxo de snapshot apenas uma vez. Quando True, o fluxo é pulado durante atualizações incrementais subsequentes após ele ser confirmado com sucesso. Uma atualização completa do alvo refaz o fluxo. Use essa opção para um carregamento de snapshot único ou backfill. O padrão é False. |
Notes
Para alvos SCD Tipo 1 em pipelines acionados, um create_auto_cdc_from_snapshot_flow() fluxo pode compartilhar um alvo com um ou mais AUTO CDC fluxos definidos em Python ou SQL. Os seguintes requisitos são aplicados:
- Dê a cada
AUTO CDCfluxo um nome único. - Use o mesmo número de chaves na mesma ordem para todos os fluxos. Nomes de chaves snapshot-flow são comparados de forma insensível a maiúsculas minúsculas com
AUTO CDCnomes de chaves. MúltiplosAUTO CDCfluxos devem usar nomes-chave e carcaças idênticos. - Use exatamente o mesmo tipo de dado para a versão snapshot e a coluna de sequenciamento de cada
AUTO CDCfluxo. - Não defina expectativas sobre o
AUTO CDC FROM SNAPSHOTfluxo. - Não use
IGNORE NULL UPDATES. Em Python, não definaignore_null_updates,ignore_null_updates_column_list, nemignore_null_updates_except_column_listnoscreate_auto_cdc_flow()fluxos. - Se você adicionar um
AUTO CDCfluxo a um destino existenteAUTO CDC FROM SNAPSHOT, o alvo não deve conter colunas de usuário cujos nomes conflitam com colunas reservadasAUTO CDCdo sistema. - Não use esse padrão para alvos SCD Tipo 2 ou bitemporais.
Para um exemplo, veja Adicionar um preenchimento a uma tabela AUTO CDC SCD Tipo 1.
Implementar o source argumento
A create_auto_cdc_from_snapshot_flow() função inclui o source argumento. Para processar instantâneos históricos, espera-se que o source argumento seja uma função lambda do Python que retorna dois valores para a create_auto_cdc_from_snapshot_flow() função: um DataFrame do Python que contém os dados de instantâneo a serem processados e uma versão de instantâneo.
A primeira invocação da função deve retornar um DataFrame e versão snapshot. Se a primeira invocação retornar None, a atualização do pipeline falha. Após pelo menos um snapshot ser processado, retorne None para indicar que não há snapshots adicionais disponíveis.
Veja a seguir a assinatura da função lambda:
lambda Any => Optional[(DataFrame, Any)]
- O argumento para a função lambda é a versão de instantâneo mais recentemente processada.
- O valor retornado da função lambda é
Noneou uma tupla de dois valores: o primeiro valor da tupla é um DataFrame que contém o instantâneo a ser processado. O segundo valor da tupla é a versão do instantâneo que representa a ordem lógica do instantâneo.
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 de pipelines do Lakeflow executa as seguintes etapas sempre que o pipeline que contém a create_auto_cdc_from_snapshot_flow() função é disparado:
- Executa a função
next_snapshot_and_versionpara carregar o próximo DataFrame de instantâneo e a versão do instantâneo correspondente. - Se a primeira invocação não devolver DataFrame, a atualização falha. Se uma invocação posterior não retorna DataFrame, a execução termina e a atualização do pipeline é marcada como completa.
- Detecta as alterações no novo instantâneo e as aplica incrementalmente à tabela de destino.
- Retorna à etapa número 1 para carregar o próximo instantâneo e sua versão.