Pular para o conteúdo principal

APIs AUTO CDC : Simplifique a captura de dados de alterações (CDC) com pipeline

Os LakeFlow Pipelines simplificam a captura de dados de alterações (CDC) com as APIs AUTO CDC e AUTO CDC FROM SNAPSHOT. Estas APIs automatizam a complexidade de computar dimensões que mudam lentamente (SCD) Tipo 1 e Tipo 2 a partir de um feed de captura de dados de alterações (CDC) ou Snapshots de banco de dados. A API AUTO CDC também suporta acompanhamento bitemporal, que registra alterações em duas dimensões de tempo (Beta). Para saber mais sobre SCD Tipo 1 e Tipo 2, consulte captura de dados de alterações (CDC) e Snapshots. Para saber mais sobre acompanhamento bitemporal, consulte CDC AUTO Bitemporal.

nota

As APIs AUTO CDC substituem as APIs APPLY CHANGES e têm a mesma sintaxe. As APIs APPLY CHANGES ainda estão disponíveis, mas a Databricks recomenda o uso das APIs AUTO CDC em seu lugar.

A API que você utiliza depende da origem dos seus dados de alteração:

  • AUTO CDC : Use isso quando o banco de dados de origem tiver um feed CDC ativado. AUTO CDC processa alterações de um feed de dados de alteração (CDF). É compatível com as interfaces SQL e Python do pipeline.
  • AUTO CDC FROM SNAPSHOT : Use isto quando o CDC não estiver habilitado no banco de dados de origem e somente Snapshots estiverem disponíveis. Esta API compara snapshots para determinar as mudanças e, em seguida, os processa. Ele é compatível com as interfaces de SQL e Python do pipeline.

Para destinos do SCD Tipo 1, você também pode combinar um fluxo AUTO CDC FROM SNAPSHOT com um ou mais fluxos AUTO CDC. Use este padrão quando uma fonte fornecer Snapshot autoritativos periódicos e um feed de CDC de menor latência entre os Snapshot.

Ambas as APIs suportam a atualização de tabelas usando SCD Tipo 1 e Tipo 2:

  • Utilize o SCD Tipo 1 para atualizar registros diretamente. não é mantido para fins de atualização de registros.
  • Use o SCD tipo 2 para manter um histórico de registros, em todas as atualizações ou em atualizações de um conjunto específico de colunas.

Somente para AUTO CDC, você também pode usar o armazenamento bitemporal, que estende o histórico do SCD Tipo 2 para rastrear alterações em duas dimensões de tempo: tempo de negócios e tempo de sistema. Bitemporal está em Beta. Consulte CDC AUTO bitemporal.

AUTO CDC também oferece suporte a atualizações parciais, onde um registro de alteração atualiza apenas um subconjunto de colunas. Consulte Aplicar atualizações parciais.

As APIs AUTO CDC não são suportadas pelo pipeline declarativo Apache Spark .

Para sintaxe e outras referências, consulte AUTO CDC INTO (pipeline), create_auto_cdc_flow e create_auto_cdc_from_snapshot_flow.

nota

Use estas APIs para atualizar tabelas em seus pipelines com base nas alterações nos dados de origem. Para saber como registrar e fazer query de informações de alteração em nível de linha para tabelas Delta, consulte Usar o feed de dados de alteração no Databricks.

Requisitos​

Para usar as APIs CDC, seu pipeline deve ser configurado para usar Serverless LakeFlow Pipelines ou as edições Pro ou Advanced do LakeFlow Pipelines.

Como funciona o AUTO CDC​

Para realizar o processamento CDC com AUTO CDC, crie uma tabela de transmissão e use a instrução AUTO CDC ... INTO em SQL ou a função create_auto_cdc_flow() em Python para especificar a origem, a chave e a sequência para o feed de alterações. Para obter uma explicação de como funcionam o sequenciamento e a lógica SCD , consulte captura de dados de alterações (CDC) e Snapshot. Veja os exemplos do AUTO CDC.

Para hidratação inicial a partir de uma fonte com alimentação de mudança, use AUTO CDC com um fluxo once e continue processando a alimentação de mudança. Consulte Replicar uma tabela RDBMS externa usando o AUTO CDC.

Para detalhes de sintaxe, consulte AUTO CDC INTO (pipeline) ou create_auto_cdc_flow.

Como funciona o AUTO CDC do Snapshot​

AUTO CDC FROM SNAPSHOT determines changes in source data by comparing in-order Snapshot. You can read snapshots from a Delta table, cloud storage files, or JDBC directly. It is supported in both the pipeline SQL and Python interfaces.

Para executar o processamento CDC com AUTO CDC FROM SNAPSHOT, crie uma tabela de transmissão e, em seguida, defina o fluxo:

  • In SQL, use the AUTO CDC ... FROM SNAPSHOT form of the CREATE FLOW statement or the embedded FLOW AUTO CDC clause of CREATE STREAMING TABLE. The source is a required FROM SNAPSHOT (snapshot_query) clause and an optional WITH VERSION (version_query) clause that selects the next snapshot version to process. See CREATE FLOW (pipelines).
  • No Python, use a função create_auto_cdc_from_snapshot_flow() para especificar o Snapshot, as chaves (key) e outros argumentos. Consulte create_auto_cdc_from_snapshot_flow.

Para obter detalhes sobre os dois padrões de ingestão e quando usar cada um, consulte Padrões de processamento de snapshot. Consulte os exemplos de AUTO CDC FROM SNAPSHOT.

Combinar o AUTO CDC FROM SNAPSHOT e o AUTO CDC​

A versão do snapshot e a coluna de sequenciamento de CDC formam um domínio de ordenação. Um evento de CDC mais recente do que o último snapshot tem precedência sobre o snapshot. Um snapshot mais recente tem precedência sobre eventos de CDC mais antigos e trata as chaves ausentes no snapshot como excluídas na versão do snapshot. Eventos de CDC mais recentes podem preservar ou restaurar essas key.

Um destino unificado deve atender aos seguintes requisitos:

  • The target uses SCD Type 1.
  • O destino tem exatamente um fluxo AUTO CDC FROM SNAPSHOT e um ou mais fluxos AUTO CDC.
  • Cada fluxo AUTO CDC tem um nome exclusivo.
  • Todos os fluxos usam o mesmo número de keys na mesma ordem. Os nomes de key do fluxo de snapshot são comparados sem diferenciação de maiúsculas e minúsculas com os nomes de key de AUTO CDC. Vários fluxos AUTO CDC devem usar nomes de key e maiúsculas e minúsculas idênticos.
  • A versão do snapshot e cada coluna de sequenciamento do CDC têm exatamente o mesmo tipo de dados.
  • O fluxo AUTO CDC FROM SNAPSHOT não define expectativas.
  • The AUTO CDC flows do not use IGNORE NULL UPDATES or its column-list variants.
  • O pipeline usa o modo Trigger.
  • If you add an AUTO CDC flow to an existing AUTO CDC FROM SNAPSHOT target, the target must not contain user columns whose names conflict with reserved AUTO CDC system columns.

Os fluxos AUTO CDC e AUTO CDC FROM SNAPSHOT podem usar a interface de pipeline em Python ou SQL. Você pode misturar fluxos em SQL e Python no mesmo destino.

Use once=True em Python ou ONCE em SQL no fluxo AUTO CDC FROM SNAPSHOT para preencher um snapshot uma vez, enquanto os fluxos AUTO CDC continuam processando novos eventos em atualizações subsequentes. Uma refresh completa do destino reexecuta o fluxo de Snapshot de uso único. Para ver um exemplo completo, consulte Adicionar um preenchimento retroativo a uma tabela AUTO CDC SCD Tipo 1.

This pattern does not support SCD Type 2 or bitemporal targets.

Use várias colunas para sequenciamento​

Para sequenciar por várias colunas (por exemplo, um carimbo de data/hora e um ID para desempatar), use STRUCT para combiná-las. A API ordena os resultados pelo primeiro campo e, em caso de empate, considera o segundo campo, e assim por diante.

SQL
SEQUENCE BY STRUCT(timestamp_col, id_col)

Exemplos de AUTO CDC​

Os exemplos a seguir demonstram o processamento de SCD Tipo 1 e Tipo 2 usando uma fonte de dados de alteração. Os dados de exemplo criam novos registros de usuário, excluem registros de usuário e atualizam registros de usuário. No exemplo SCD Tipo 1, as últimas operações UPDATE chegam atrasadas e são descartadas da tabela de destino, demonstrando o tratamento de eventos fora de ordem.

A seguir, estão os registros de entrada usados nestes exemplos. Esses dados são criados executando a consulta na seção Criar dados de exemplo .

userId

name

city

operation

sequenceNum

124

Raul

Oaxaca

INSERT

1

123

Isabel

Monterrey

INSERT

1

125

Mercedes

Tijuana

INSERT

2

126

Lily

Cancun

INSERT

2

123

null

null

DELETE

6

125

Mercedes

Guadalajara

UPDATE

6

125

Mercedes

Mexicali

UPDATE

5

123

Isabel

Chihuahua

UPDATE

5

userId

name

city

operation

sequenceNum

124

Raul

Oaxaca

INSERT

1

123

Isabel

Monterrey

INSERT

1

125

Mercedes

Tijuana

INSERT

2

126

Lily

Cancun

INSERT

2

123

null

null

DELETE

6

125

Mercedes

Guadalajara

UPDATE

6

125

Mercedes

Mexicali

UPDATE

5

123

Isabel

Chihuahua

UPDATE

5

Se você remover o comentário da última linha na consulta de geração de dados de exemplo, será inserido o seguinte registro que especifica o truncamento da tabela (limpeza da tabela) em sequenceNum=3:

userId

name

city

operation

sequenceNum

null

null

null

TRUNCATE

3

userId

name

city

operation

sequenceNum

null

null

null

TRUNCATE

3

nota

Todos os exemplos a seguir incluem opções para especificar as operações DELETE e TRUNCATE , mas cada uma é opcional.

Criar dados de exemplo​

Execute as seguintes instruções para criar um dataset de exemplo. Este código não se destina a ser executado como parte de uma definição pipeline . execute-o a partir da pasta de exploração do seu pipeline, em vez da pasta de transformações.

SQL
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;

CREATE TABLE main.cdc_tutorial.users_cdf
AS SELECT
col1 AS userId,
col2 AS name,
col3 AS city,
col4 AS operation,
col5 AS sequenceNum
FROM (
VALUES
-- Initial load.
(124, "Raul", "Oaxaca", "INSERT", 1),
(123, "Isabel", "Monterrey", "INSERT", 1),
-- New users.
(125, "Mercedes", "Tijuana", "INSERT", 2),
(126, "Lily", "Cancun", "INSERT", 2),
-- Isabel is removed from the system and Mercedes moved to Guadalajara.
(123, null, null, "DELETE", 6),
(125, "Mercedes", "Guadalajara", "UPDATE", 6),
-- This batch of updates arrived out of order. The batch at sequenceNum 6 is the final state.
(125, "Mercedes", "Mexicali", "UPDATE", 5),
(123, "Isabel", "Chihuahua", "UPDATE", 5)
-- Uncomment to test TRUNCATE.
-- ,(null, null, null, "TRUNCATE", 3)
);

Atualizações do Processo SCD Tipo 1​

O SCD Type 1 mantém apenas a versão mais recente de cada registro. O exemplo a seguir faz a leitura do feed de dados de alteração criado acima e aplica as alterações a um destino de tabela de streaming. Você precisa de um pipeline para executar este código. Consulte Develop and debug ETL pipelines with the Lakeflow Pipelines Editor.

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr

@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")

dp.create_streaming_table("users_current")

dp.create_auto_cdc_flow(
target = "users_current",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)

Após executar o exemplo do SCD tipo 1, a tabela de destino contém os seguintes registros:

userId

name

city

124

Raul

Oaxaca

125

Mercedes

Guadalajara

126

Lily

Cancun

userId

name

city

124

Raul

Oaxaca

125

Mercedes

Guadalajara

126

Lily

Cancun

O usuário 123 (Isabel) foi excluído e não aparece mais. O usuário 125 (Mercedes) mostra apenas a cidade mais recente (Guadalajara) porque o SCD Tipo 1 sobrescreve os valores anteriores. O UPDATE anterior em sequenceNum=5 foi removido porque uma atualização posterior em sequenceNum=6 chegou.

Após executar o exemplo com o registro TRUNCATE descomentado, a tabela é limpa em sequenceNum=3. Isso significa que os registros 124 e 126 não estão na tabela, e a tabela de destino final contém apenas o seguinte registro:

userId

name

city

125

Mercedes

Guadalajara

userId

name

city

125

Mercedes

Guadalajara

Atualizações do Processo SCD Tipo 2​

SCD Tipo 2 preserva um histórico completo de alterações, criando novas linhas para cada versão de um registro, com colunas __START_AT e __END_AT indicando quando cada versão estava ativa.

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr

@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")

dp.create_streaming_table("users_history")

dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2"
)

Após executar o exemplo do SCD tipo 2, a tabela de destino contém os seguintes registros:

userId

name

city

__START_AT

__END_AT

123

Isabel

Monterrey

1

5

123

Isabel

Chihuahua

5

6

124

Raul

Oaxaca

1

null

125

Mercedes

Tijuana

2

5

125

Mercedes

Mexicali

5

6

125

Mercedes

Guadalajara

6

null

126

Lily

Cancun

2

null

userId

name

city

__START_AT

__END_AT

123

Isabel

Monterrey

1

5

123

Isabel

Chihuahua

5

6

124

Raul

Oaxaca

1

null

125

Mercedes

Tijuana

2

5

125

Mercedes

Mexicali

5

6

125

Mercedes

Guadalajara

6

null

126

Lily

Cancun

2

null

A mesa preserva toda a história. O usuário 123 possui duas versões (terminada na sequência 6 quando excluída). O usuário 125 possui três versões que mostram alterações na cidade. Registros com __END_AT = null estão atualmente ativos.

Rastreie um subconjunto de colunas com SCD Tipo 2.​

Por default, SCD Tipo 2 cria uma nova versão sempre que o valor de qualquer coluna é alterado. Você pode especificar um subconjunto de colunas para rastrear, de forma que as alterações em outras colunas atualizem a versão atual no mesmo local, em vez de gerar um novo registro no histórico.

O exemplo a seguir exclui a coluna city da história acompanhamento:

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr

@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")

dp.create_streaming_table("users_history")

dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2",
track_history_except_column_list = ["city"]
)

Como as alterações city não são rastreadas, as atualizações de cidade sobrescrevem a linha atual em vez de criar uma nova versão. A tabela de destino contém os seguintes registros:

userId

name

city

__START_AT

__END_AT

123

Isabel

Chihuahua

1

6

124

Raul

Oaxaca

1

null

125

Mercedes

Guadalajara

2

null

126

Lily

Cancun

2

null

userId

name

city

__START_AT

__END_AT

123

Isabel

Chihuahua

1

6

124

Raul

Oaxaca

1

null

125

Mercedes

Guadalajara

2

null

126

Lily

Cancun

2

null

CDC AUTOMÁTICO DE exemplos de instantâneos​

As seções a seguir fornecem exemplos de uso de AUTO CDC FROM SNAPSHOT para processar Snapshot em tabelas de destino SCD Tipo 1 ou Tipo 2. Para obter informações sobre quando usar esta API, consulte captura de dados de alterações (CDC) e Snapshot.

Os exemplos nesta seção incluem implementações em Python e SQL. Para ver a sintaxe SQL completa, consulte CREATE FLOW (pipelines).

Example: Process periodic Snapshot​

Use this approach when Snapshot arrive regularly and in order. The Python interface uses the pipeline update order for versioning. For SQL, the example uses a separate version table to identify when a new snapshot is available.

Você pode ler Snapshots de vários tipos de origem, incluindo tabelas Delta , arquivos de armazenamento cloud e conexões JDBC .

o passo 1: Criar dados de exemplo​

Crie uma tabela contendo dados de instantâneo. Execute o seguinte código a partir de um Notebook ou Databricks SQL na pasta explorations do seu pipeline:

SQL
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;

CREATE TABLE main.cdc_tutorial.snapshot (
userId INT,
city STRING
);

CREATE TABLE main.cdc_tutorial.snapshot_version (
version BIGINT
);

INSERT INTO main.cdc_tutorial.snapshot VALUES
(1, 'Oaxaca'),
(2, 'Monterrey'),
(3, 'Tijuana');

INSERT INTO main.cdc_tutorial.snapshot_version VALUES (0);

o passo 2: execução AUTO CDC FROM Snapshot​

Você precisa de um pipeline para executar o código nesta etapa. Consulte Develop and debug ETL pipelines with the Lakeflow Pipelines Editor.

Para Python, escolha um tipo de origem para a view de snapshot (o código de criação de exemplo gera uma tabela Delta):

Opção A: Ler de uma tabela Delta

Python
from pyspark import pipelines as dp

@dp.view(name="source")
def source():
return spark.read.table("main.cdc_tutorial.snapshot")

Opção B: Ler do armazenamento cloud

Python
from pyspark import pipelines as dp

@dp.view(name="source")
def source():
return spark.read.format("csv").option("header", True).load("<snapshot-path>")

Opção C: Ler do JDBC (somente compute clássica)

Python
from pyspark import pipelines as dp

@dp.view(name="source")
def source():
return (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)

Then add the target table and flow. The SQL implementation uses the snapshot_version table to identify when a new snapshot is available:

Python
dp.create_streaming_table("target")

dp.create_auto_cdc_from_snapshot_flow(
target = "target",
source = "source",
keys = ["userId"],
stored_as_scd_type = 2
)

Após a primeira execução do pipeline, todos os registros são inseridos como linhas ativas:

userId

city

__START_AT

__END_AT

1

Oaxaca

0

null

2

Monterrey

0

null

3

Tijuana

0

null

userId

city

__START_AT

__END_AT

1

Oaxaca

0

null

2

Monterrey

0

null

3

Tijuana

0

null

nota

Para usar o SCD Tipo 1 e manter apenas o estado atual, defina stored_as_scd_type=1 em Python ou STORED AS SCD TYPE 1 em SQL. Nesse caso, a tabela de destino não inclui as colunas __START_AT e __END_AT.

o passo 3: Simule um novo Snapshot e execute novamente​

Atualize a tabela de origem para simular a chegada de um novo snapshot (execute este código de um notebook ou arquivo SQL na pasta explorations do seu pipeline):

SQL
TRUNCATE TABLE main.cdc_tutorial.snapshot;

INSERT INTO main.cdc_tutorial.snapshot VALUES
(2, 'Carmel'),
(3, 'Los Angeles'),
(4, 'Death Valley'),
(6, 'Kings Canyon');

UPDATE main.cdc_tutorial.snapshot_version SET version = 1;

executar o pipeline novamente. AUTO CDC FROM SNAPSHOT compara o novo Snapshot com o anterior e detecta que o usuário 1 foi excluído, os usuários 2 e 3 foram atualizados e os usuários 4 e 6 foram inseridos. Isso gera um feed de alterações e usa AUTO CDC para criar a tabela de saída.

Após a segunda execução com SCD Tipo 2, a tabela de destino contém os seguintes registros:

userId

city

__START_AT

__END_AT

1

Oaxaca

0

1

2

Monterrey

0

1

2

Carmelo

1

null

3

Tijuana

0

1

3

Los Angeles

1

null

4

Vale da Morte

1

null

6

Kings Canyon

1

null

userId

city

__START_AT

__END_AT

1

Oaxaca

0

1

2

Monterrey

0

1

2

Carmelo

1

null

3

Tijuana

0

1

3

Los Angeles

1

null

4

Vale da Morte

1

null

6

Kings Canyon

1

null

O usuário 1 foi encerrado (excluído). Os usuários 2 e 3 têm duas versões cada, mostrando as alterações em suas cidades. Os usuários 4 e 6 foram inseridos recentemente.

Após a segunda execução com SCD Tipo 1, a tabela de destino mostra apenas o estado atual:

userId

city

2

Carmelo

3

Los Angeles

4

Vale da Morte

6

Kings Canyon

userId

city

2

Carmelo

3

Los Angeles

4

Vale da Morte

6

Kings Canyon

Exemplo: Captura instantânea do processo usando funções de versão​

Use esta abordagem quando precisar de controle explícito sobre a ordenação dos Snapshot. Por exemplo, use esta abordagem quando vários snapshots chegarem ao mesmo tempo ou chegarem fora de ordem. Você define como selecionar o próximo snapshot e seu número de versão. A API processa os snapshots em ordem crescente de versão:

  • Se vários snapshots estiverem no armazenamento, o seletor de exemplo processará os snapshots disponíveis na ordem de versão.
  • Se um snapshot chegar fora de ordem (por exemplo, snapshot_3 chegar após snapshot_4), o seletor de exemplo o exclui porque sua versão não é superior à última versão confirmada.
  • If there are no new snapshots, the version function or query returns no result and no processing occurs.

o passo 1: Preparar arquivos de instantâneo​

Crie arquivos CSV contendo dados de Snapshot e adicione-os a um volume ou local de armazenamento cloud . Nomeie os arquivos cronologicamente (por exemplo, snapshot_1.csv, snapshot_2.csv).

Cada arquivo deve conter colunas para userId e city. Por exemplo:

snapshot_1.csv :

userId

city

1

Oaxaca

2

Monterrey

3

Tijuana

userId

city

1

Oaxaca

2

Monterrey

3

Tijuana

snapshot_2.csv :

userId

city

2

Carmelo

3

Los Angeles

4

Vale da Morte

userId

city

2

Carmelo

3

Los Angeles

4

Vale da Morte

o passo 2: execução AUTO CDC FROM Snapshot com função de versão​

Em um pipeline, crie um novo arquivo Python ou SQL na pasta transformations e adicione o seguinte código. Em seguida, execute o pipeline. Consulte Desenvolver e depurar pipelines de ETL com o Editor de Pipelines LakeFlow Pipelines.

Python
from pyspark import pipelines as dp
from typing import Optional, Tuple
from pyspark.sql import DataFrame

def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Optional[Tuple[DataFrame, int]]:
snapshot_dir = "/Volumes/main/cdc_tutorial/snapshots/" # or the location you created the sample data

files = dbutils.fs.ls(snapshot_dir)
snapshot_files = [f.name for f in files if f.name.startswith("snapshot_") and f.name.endswith(".csv")]

snapshot_versions = []
for filename in snapshot_files:
try:
version = int(filename.replace("snapshot_", "").replace(".csv", ""))
snapshot_versions.append(version)
except ValueError:
continue

snapshot_versions.sort()

if latest_snapshot_version is None:
if snapshot_versions:
next_version = snapshot_versions[0]
else:
return None
else:
next_versions = [v for v in snapshot_versions if v > latest_snapshot_version]
if next_versions:
next_version = next_versions[0]
else:
return None

snapshot_path = f"{snapshot_dir}snapshot_{next_version}.csv"
df = spark.read.format("csv").option("header", True).load(snapshot_path)
return (df, next_version)


dp.create_streaming_table("main.cdc_tutorial.target_versioned")

dp.create_auto_cdc_from_snapshot_flow(
target = "main.cdc_tutorial.target_versioned",
source = next_snapshot_and_version,
keys = ["userId"],
stored_as_scd_type = 2
)
nota

Para usar o SCD Tipo 1, defina stored_as_scd_type=1 em Python ou STORED AS SCD TYPE 1 em SQL.

Após o processamento snapshot_1.csv, a tabela de destino contém os seguintes registros:

userId

city

__START_AT

__END_AT

1

Oaxaca

1

null

2

Monterrey

1

null

3

Tijuana

1

null

userId

city

__START_AT

__END_AT

1

Oaxaca

1

null

2

Monterrey

1

null

3

Tijuana

1

null

Após o processamento snapshot_2.csv, a tabela de destino contém os seguintes registros:

userId

city

__START_AT

__END_AT

1

Oaxaca

1

2

2

Monterrey

1

2

2

Carmelo

2

null

3

Tijuana

1

2

3

Los Angeles

2

null

4

Vale da Morte

2

null

userId

city

__START_AT

__END_AT

1

Oaxaca

1

2

2

Monterrey

1

2

2

Carmelo

2

null

3

Tijuana

1

2

3

Los Angeles

2

null

4

Vale da Morte

2

null

nota

Lembre-se de que, para SCD Tipo 1, a tabela é exatamente igual ao Snapshot mais recente. A diferença é que as consultas subsequentes podem usar o feed de alterações para processar apenas os registros modificados.

o passo 3: Adicionar novo instantâneo​

Adicione um novo arquivo CSV ao local de armazenamento com os dados modificados (por exemplo, valores de cidade alterados, novas linhas ou linhas removidas). Em seguida, execute o pipeline novamente para processar o novo Snapshot.

Limitações​

Recurso adicional​