Pular para o conteúdo principal

Tópicos avançados do AUTO CDC

Além das APIs AUTO CDC e AUTO CDC FROM SNAPSHOT básicas, você pode executar DML em tabelas de destino, ler fluxos de dados de alterações de destinos CDC, monitorar métricas de processamento, aplicar atualizações parciais e rastrear alterações com armazenamento bitemporal. Para uma introdução às APIs AUTO CDC, consulte As APIs AUTO CDC: Simplifique a captura de dados de alterações com pipelines.

Adicionar, alterar ou excluir dados em uma tabela de transmissão de destino.

Se o seu pipeline publicar tabelas no Unity Catalog, você poderá usar instruções de linguagem de manipulação de dados (DML), incluindo instruções insert, update, delete e merge , para modificar as tabelas de transmissão de destino criadas pelas instruções AUTO CDC ... INTO .

nota
  • As declarações DML que modificam o esquema de tabela de uma tabela de streaming não são suportadas. Certifique-se de que suas instruções DML não tentem evoluir o esquema da tabela.
  • Instruções DML que atualizam uma tabela de transmissão só podem ser executadas em um cluster Unity Catalog compartilhado ou em um SQL warehouse usando Databricks Runtime 13.3 LTS ou superior.
  • Como a transmissão requer uma fonte de dados somente para anexação, se o seu processamento exigir transmissão de uma tabela de transmissão de origem com alterações (por exemplo, por meio de instruções DML), defina o sinalizador skipChangeCommits ao ler a tabela de transmissão de origem. Quando skipChangeCommits está definido, as transações que excluem ou modificam registros na tabela de origem são ignoradas. Se o seu processamento não exigir uma tabela de transmissão, você pode usar uma view materializada (que não possui a restrição de somente acréscimo) como tabela de destino.

Como o pipeline usa uma coluna SEQUENCE BY especificada e propaga valores de sequenciamento apropriados para as colunas __START_AT e __END_AT da tabela de destino (para SCD Tipo 2), você deve garantir que as instruções DML usem valores válidos para essas colunas para manter a ordem adequada dos registros. Consulte Como o AUTO CDC funciona.

Para obter mais informações sobre o uso de instruções DML com tabelas de transmissão, consulte Adicionar, alterar ou excluir dados em uma tabela de transmissão.

O exemplo seguinte insere um registro ativo com uma sequência inicial de 5:

SQL
INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);
dica

Se você precisar renomear as colunas __START_AT e __END_AT na sua tabela de destino SCD Tipo 2 (por exemplo, para corresponder aos requisitos do esquema subsequente), crie uma view sobre a tabela de destino:

SQL
CREATE VIEW my_employees_view AS
SELECT
*,
__START_AT AS valid_from,
__END_AT AS valid_to
FROM my_scd2_target_table;

Leia um feed de dados de alterações de uma tabela de metas AUTO CDC.

No Databricks Runtime 15.2 e versões superiores, você pode ler um feed de dados de alteração de uma tabela de transmissão que é o alvo de consultas AUTO CDC ou AUTO CDC FROM SNAPSHOT da mesma forma que você lê um feed de dados de alteração de outras tabelas Delta . Os seguintes itens são necessários para ler o fluxo de dados de alterações de uma tabela de transmissão de destino:

  • A tabela de transmissão de destino deve ser publicada no Unity Catalog. Consulte Usar Unity Catalog com o pipeline.
  • Para ler o fluxo de dados de alterações da tabela de transmissão de destino, você precisa usar Databricks Runtime 15.2 ou superior. Para ler o fluxo de dados de alterações em um pipeline diferente, o pipeline deve ser configurado para usar Databricks Runtime 15.2 ou superior.

Você lê o feed de dados de alteração de uma tabela de transmissão de destino que foi criada em um LakeFlow Pipelines da mesma forma que lê um feed de dados de alteração de outras tabelas Delta. Para saber mais sobre o uso da funcionalidade de feed de dados de alteração Delta, incluindo exemplos em Python e SQL, consulte Usar feed de dados de alteração no Databricks.

nota

O registro do feed de dados de alteração inclui metadados que identificam o tipo de evento de alteração. Quando um registro é atualizado em uma tabela, os metadados para os registros de alteração associados normalmente incluem valores _change_type definidos para eventos update_preimage e update_postimage .

No entanto, os valores _change_type são diferentes se forem feitas atualizações na tabela de transmissão de destino que incluam a alteração dos valores key primária. Quando as alterações incluem atualizações na chave primária, os campos de metadados _change_type são definidos para os eventos insert e delete . As alterações na chave primária podem ocorrer quando atualizações manuais são feitas em um dos campos key com uma instrução UPDATE ou MERGE ou, para tabelas SCD tipo 2, quando o campo __start_at muda para refletir um valor de sequência de início anterior.

A consulta AUTO CDC determina os valores key primária, que diferem para o processamento SCD tipo 1 e SCD tipo 2:

SCD Type

Primary key

SCD type 1, and the pipelines Python interface

The primary key is the value of the keys parameter in the create_auto_cdc_flow() function. For the SQL interface the primary key is the columns defined by the KEYS clause in the AUTO CDC ... INTO statement.

SCD type 2

The primary key is the keys parameter or KEYS clause plus the return value from the coalesce(__START_AT, __END_AT) operation, where __START_AT and __END_AT are the corresponding columns from the target streaming table. This uses __START_AT when available, and __END_AT when __START_AT is null (for example, the initial record).

SCD Type

Primary key

SCD type 1, and the pipelines Python interface

The primary key is the value of the keys parameter in the create_auto_cdc_flow() function. For the SQL interface the primary key is the columns defined by the KEYS clause in the AUTO CDC ... INTO statement.

SCD type 2

The primary key is the keys parameter or KEYS clause plus the return value from the coalesce(__START_AT, __END_AT) operation, where __START_AT and __END_AT are the corresponding columns from the target streaming table. This uses __START_AT when available, and __END_AT when __START_AT is null (for example, the initial record).

Obtenha dados sobre registros processados por uma consulta CDC em andamento.

nota

As seguintes métricas são capturadas apenas por consultas AUTO CDC e não por consultas AUTO CDC FROM SNAPSHOT .

As seguintes métricas são capturadas por consultas AUTO CDC :

  • num_upserted_rows : O número de linhas de saída inseridas no dataset durante uma atualização.
  • num_deleted_rows : O número de linhas de saída existentes excluídas do dataset durante uma atualização.

As métricas num_output_rows , saída para fluxos não-CDC, não são capturadas para consultas AUTO CDC .

Aplicar atualizações parciais

Quando uma fonte envia apenas as colunas que foram alteradas, AUTO CDC deve distinguir entre uma coluna que está ausente de um registro de alteração, que deve deixar o valor de destino inalterado, e uma coluna que é definida explicitamente como null, que deve substituir o valor de destino por null. By default, IGNORE NULL UPDATES trata cada null como um marcador de "não atualizar", portanto, não pode aplicar um null explícito. Para resolver essa ambiguidade, escolha um dos três métodos a seguir:

Método

Quando usar

Comportamento

IGNORE NULL UPDATES ON columnList

Um pequeno conjunto fixo de colunas deve ignorar os valores null, enquanto todas as outras colunas aplicam valores explícitos null.

As colunas listadas mantêm seu valor de destino existente quando o valor de entrada é null. Todas as outras colunas aplicam valores null explícitos.

IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList)

A maioria das colunas deve ignorar os valores null, e apenas algumas devem aplicar valores null explícitos.

As colunas listadas aplicam valores explícitos de null. Todas as outras colunas mantêm seu valor de destino existente quando o valor de entrada é null.

COLUMNS TO UPDATE

Cada registro de alteração atualiza um conjunto diferente de colunas, ou o conjunto de colunas atualizáveis muda ao longo do tempo.

Uma coluna de origem nomeia as colunas a serem atualizadas para cada registro de alteração. As colunas listadas são gravadas da origem, incluindo valores null explícitos. As colunas que não estão listadas mantêm o valor de destino existente.

Método

Quando usar

Comportamento

IGNORE NULL UPDATES ON columnList

Um pequeno conjunto fixo de colunas deve ignorar os valores null, enquanto todas as outras colunas aplicam valores explícitos null.

As colunas listadas mantêm seu valor de destino existente quando o valor de entrada é null. Todas as outras colunas aplicam valores null explícitos.

IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList)

A maioria das colunas deve ignorar os valores null, e apenas algumas devem aplicar valores null explícitos.

As colunas listadas aplicam valores explícitos de null. Todas as outras colunas mantêm seu valor de destino existente quando o valor de entrada é null.

COLUMNS TO UPDATE

Cada registro de alteração atualiza um conjunto diferente de colunas, ou o conjunto de colunas atualizáveis muda ao longo do tempo.

Uma coluna de origem nomeia as colunas a serem atualizadas para cada registro de alteração. As colunas listadas são gravadas da origem, incluindo valores null explícitos. As colunas que não estão listadas mantêm o valor de destino existente.

COLUMNS TO UPDATE Não pode ser combinado com IGNORE NULL UPDATES e não é suportado para tabelas bitemporais.

Como regra geral, escolha COLUMNS TO UPDATE quando o produtor souber quais colunas foram alteradas em cada registro e puder transportar essa informação em uma coluna de origem, como quando vários produtores gravam na mesma fonte ou o conjunto de colunas atualizáveis cresce ao longo do tempo. Escolha IGNORE NULL UPDATES ON quando o proprietário do pipeline souber antecipadamente o conjunto fixo de colunas atualizáveis e preferir controlá-las no código do pipeline.

O exemplo a seguir usa uma coluna de origem chamada columnsToUpdate para controlar quais colunas cada registro de alteração atualiza, incluindo colunas definidas explicitamente como null:

Python
from pyspark import pipelines as dp

dp.create_streaming_table("target")

dp.create_auto_cdc_flow(
target = "target",
source = "cdc_source",
keys = ["id"],
sequence_by = "sequenceNum",
stored_as_scd_type = 1,
columns_to_update = "columnsToUpdate"
)

Para a referência completa de parâmetros, consulte AUTO CDC INTO (pipelines) e create_auto_cdc_flow.

AUTO CDC Bitemporal

info

Beta

O AUTO CDC bitemporal está em Beta.

SCD Tipo 1 e Tipo 2 são unitemporais: elas rastreiam mudanças em uma única dimensão de tempo. Bitemporal estende o histórico do SCD Tipo 2 para rastrear mudanças em duas dimensões de tempo e distinguir entre duas perspectivas:

  • **Tempo de negócio**: quando o evento realmente aconteceu.
  • Hora do sistema : quando o sistema registrou ou ingeriu o evento.

Assim como o SCD Tipo 2, o bitemporal preserva um histórico completo de registros. Ele adiciona uma segunda linha do tempo para que seja possível reconstruir tanto o que os dados mostraram quanto o que o sistema acreditava em qualquer ponto no passado.

Por exemplo, um fundo de hedge ingere dados de ações de um sistema de origem. O preço das ações da Acme Corp. muda em 1º de janeiro, mas o fundo não ingere essa atualização até 5 de janeiro. O CDC AUTOMÁTICO Bitemporal permite que o fundo responda a duas perguntas distintas: qual era o preço real das ações da Acme Corp. em 1º de janeiro (hora comercial) e qual preço o sistema acreditava quando o fundo tomou decisões de negociação em 3 de janeiro (hora do sistema). A capacidade de distinguir entre essas linhas do tempo é útil para auditoria, relatórios regulatórios e tomada de decisões financeiras.

Para habilitar o processamento bitemporal, defina STORED AS BITEMPORAL (SQL) ou stored_as_scd_type="bitemporal" (Python), use SEQUENCE BY para a coluna de tempo de negócio e use SYSTEM SEQUENCE BY para a coluna de tempo de sistema. A tabela de destino adiciona as colunas __SYSTEM_START_AT e __SYSTEM_END_AT junto com as colunas SCD Tipo 2 __START_AT e __END_AT. Para obter detalhes de sintaxe, consulte AUTO CDC INTO (pipelines) ou create_auto_cdc_flow.

Exemplos de CDC AUTO Bitemporal

O exemplo a seguir cria uma tabela de destino bitemporal a partir de um pequeno conjunto de eventos sintéticos CDC. A coluna bt contém o horário comercial e a coluna st contém o horário do sistema.

Python
from pyspark import pipelines as dp

# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")

@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
return spark.createDataFrame(
[
(1, "x10", "y10", 10, 100),
(1, "x20", "y20", 20, 200)
],
schema="id INT, x STRING, y STRING, bt INT, st INT",
)

# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")

dp.create_auto_cdc_flow(
target = "target_bitemporal",
source = "cdc_source",
keys = ["id"],
sequence_by = "bt",
system_sequence_by = "st",
stored_as_scd_type = "bitemporal"
)

A seguinte sequência de alterações mostra como uma tabela bitemporal registra uma inserção, uma atualização, uma atualização fora de ordem e uma exclusão para uma única empresa. A coluna de sequenciamento gera as colunas __START_AT e __END_AT (horário comercial), e a coluna de sequenciamento do sistema gera as colunas __SYSTEM_START_AT e __SYSTEM_END_AT (horário do sistema):

Coluna

Descrição

__START_AT

O horário comercial em que esta linha se tornou válida.

__END_AT

O horário comercial em que a validade desta linha termina. null se válido indefinidamente.

__SYSTEM_START_AT

O horário do sistema no qual os dados desta linha e o intervalo de tempo comercial são considerados verdadeiros.

__SYSTEM_END_AT

O horário do sistema no qual os dados desta linha e o intervalo de tempo comercial são considerados invalidados. null se for conhecido como verdadeiro indefinidamente.

Coluna

Descrição

__START_AT

O horário comercial em que esta linha se tornou válida.

__END_AT

O horário comercial em que a validade desta linha termina. null se válido indefinidamente.

__SYSTEM_START_AT

O horário do sistema no qual os dados desta linha e o intervalo de tempo comercial são considerados verdadeiros.

__SYSTEM_END_AT

O horário do sistema no qual os dados desta linha e o intervalo de tempo comercial são considerados invalidados. null se for conhecido como verdadeiro indefinidamente.

O sistema lida com eventos que chegam em qualquer ordem em ambas as linhas do tempo. Quando um evento chega com um tempo de negócio ou tempo de sistema anterior do que eventos já processados, o sistema corrige a história afetada em vez de apenas anexar ao final.

Alteração 1: Inserir

A Empresa A é adicionada em 18 de julho de 2025 10:01:00 (horário comercial), mas não é ingerida até 10:05:00 (horário do sistema).

Entrada:

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv1

18/07/2025 10:01:00

18/07/2025 10:05:00

INSERT

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv1

18/07/2025 10:01:00

18/07/2025 10:05:00

INSERT

Saída:

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

null

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

null

XFv1 é válido a partir das 10:01:00, sem término conhecido. O sistema soube desse fato no horário do sistema 10:05:00, sem término conhecido.

Alteração 2: Atualização

A Empresa A é atualizada em 18 de julho de 2025 12:15:43 (hora comercial), e o sistema consome o evento às 12:20:00 (hora do sistema). O sistema preserva tanto aquilo que acreditava saber antes da atualização quanto a história de negócios corrigida após a ingestão da atualização.

Entrada:

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv2

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

UPDATE

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv2

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

UPDATE

Saída:

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

18 de julho de 2025 12:20:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

null

A

XFv2

18/07/2025 12:15:43

null

18 de julho de 2025 12:20:00

null

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

18 de julho de 2025 12:20:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

null

A

XFv2

18/07/2025 12:15:43

null

18 de julho de 2025 12:20:00

null

Acreditava-se que XFv1 era válido a partir das 10:01:00 sem fim conhecido, e o sistema manteve essa crença das 10:05:00 até as 12:20:00. XFv1 é agora conhecido por ser válido apenas até 12:15:43, uma história corrigida efetiva a partir do tempo do sistema 12:20:00 sem fim conhecido. XFv2 é válido a partir das 12:15:43 sem fim conhecido, e foi aprendido no horário do sistema 12:20:00.

Alteração 3: Atualização fora de ordem

Uma atualização fora de ordem chega indicando que a Empresa A foi realmente atualizada em 18/07/2025 12:05:00 (tempo de negócio), mas não é ingerida até 12:25:00 (tempo do sistema). Quando uma atualização chega mais tarde no tempo do sistema, mas com um tempo de negócio precedente, o sistema corrige o tempo de negócio histórico e preserva tanto o que ele acreditava antes da atualização fora de ordem quanto o histórico corrigido.

Entrada:

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv3

18/07/2025 12:05:00

18 de julho de 2025 12:25:00

UPDATE

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv3

18/07/2025 12:05:00

18 de julho de 2025 12:25:00

UPDATE

Saída:

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

18 de julho de 2025 12:20:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

18 de julho de 2025 12:25:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:05:00

18 de julho de 2025 12:25:00

null

A

XFv3

18/07/2025 12:05:00

18/07/2025 12:15:43

18 de julho de 2025 12:25:00

null

A

XFv2

18/07/2025 12:15:43

null

18 de julho de 2025 12:20:00

null

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

18 de julho de 2025 12:20:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

18 de julho de 2025 12:25:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:05:00

18 de julho de 2025 12:25:00

null

A

XFv3

18/07/2025 12:05:00

18/07/2025 12:15:43

18 de julho de 2025 12:25:00

null

A

XFv2

18/07/2025 12:15:43

null

18 de julho de 2025 12:20:00

null

Acreditava-se que o XFv1 era válido das 10:01:00 às 12:15:43, e essa crença agora é válida no tempo de sistema até as 12:25:00. A nova atualização corrige a validade comercial do XFv1 para terminar às 12:05:00, um histórico corrigido efetivo a partir do tempo de sistema 12:25:00. Agora, sabe-se que o XFv3 é válido das 12:05:00 até as 12:15:43, uma crença válida no tempo de sistema a partir das 12:25:00, sem fim conhecido.

Alteração 4: Excluir

A Empresa A é excluída em 18/07/2025 12:30:00, e o sistema consome o evento às 12:30:00. Como uma operação de exclusão representa o fim da existência comercial da entidade, o sistema não cria nenhuma linha de substituição. XFv2 aparece em duas linhas, preservando um registro de auditoria completo de quando a empresa deixou de existir e quando o sistema soube da exclusão.

Entrada:

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv2

18/07/2025 12:30:00

18/07/2025 12:30:00

DELETE

CompanyId

pontos de dados

Sequenciamento

Sequenciamento do Sistema

Operação

A

XFv2

18/07/2025 12:30:00

18/07/2025 12:30:00

DELETE

Saída:

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

18 de julho de 2025 12:20:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

18 de julho de 2025 12:25:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:05:00

18 de julho de 2025 12:25:00

null

A

XFv3

18/07/2025 12:05:00

18/07/2025 12:15:43

18 de julho de 2025 12:25:00

null

A

XFv2

18/07/2025 12:15:43

null

18 de julho de 2025 12:20:00

18/07/2025 12:30:00

A

XFv2

18/07/2025 12:15:43

18/07/2025 12:30:00

18/07/2025 12:30:00

null

CompanyId

pontos de dados

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

18/07/2025 10:01:00

null

18/07/2025 10:05:00

18 de julho de 2025 12:20:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:15:43

18 de julho de 2025 12:20:00

18 de julho de 2025 12:25:00

A

XFv1

18/07/2025 10:01:00

18/07/2025 12:05:00

18 de julho de 2025 12:25:00

null

A

XFv3

18/07/2025 12:05:00

18/07/2025 12:15:43

18 de julho de 2025 12:25:00

null

A

XFv2

18/07/2025 12:15:43

null

18 de julho de 2025 12:20:00

18/07/2025 12:30:00

A

XFv2

18/07/2025 12:15:43

18/07/2025 12:30:00

18/07/2025 12:30:00

null

XFv2 era válido a partir de 12:15:43 sem fim conhecido, e o sistema manteve essa crença de 12:20:00 a 12:30:00. Após a exclusão ser ingerida, o XFv2 é conhecido por ser válido apenas até 12:30:00, uma história corrigida com efeito a partir do horário do sistema 12:30:00.

Quais objetos de dados são usados para o processamento do CDC em um pipeline?

Quando você declara a tabela de destino no Hive metastore, são criadas duas estruturas de dados:

  • Uma visualização usando o nome atribuído à tabela de destino.
  • Uma tabela de suporte interna usada pelo pipeline para gerenciar o processamento do CDC. Esta tabela é nomeada adicionando __apply_changes_storage_ antes do nome da tabela de destino.

Por exemplo, se você declarar uma tabela de destino chamada dp_cdc_target, você verá uma view chamada dp_cdc_target e uma tabela chamada __apply_changes_storage_dp_cdc_target no metastore. Consulte a view para acessar os dados processados. Não modifique a tabela de suporte diretamente.

nota

Essas estruturas de dados se aplicam apenas ao processamento AUTO CDC , não ao processamento AUTO CDC FROM SNAPSHOT . Elas também se aplicam somente ao Hive metastore, não Unity Catalog.