Pular para o conteúdo principal

Substituição parcial de Snapshot com fluxos REPLACE USING

info

Beta

Esse recurso está em Beta.

Um fluxo REPLACE USING mantém uma tabela de destino sincronizada com uma origem de transmissão: ele substitui todas as linhas que correspondem às colunas de key especificadas e deixa todos os outros dados inalterados.

Uma coluna SEQUENCE BY ordena as atualizações para que o resultado esteja correto mesmo quando as atualizações chegam fora de ordem. Para cada key, a sequência mais alta vence, e uma linha de sequência inferior nunca substitui uma superior que já esteja no destino. As linhas que compartilham a mesma key e a mesma sequência são anexadas em vez de substituídas.

Como o REPLACE USING funciona

Considere uma tabela de eventos que contém eventos de clique e conversão para duas regiões, sequenciados por seq:

id_da_região

tipo de dispositivo

event_type

seq

1

iOS

clique

1

1

Android

conversão

1

2

iOS

clique

1

2

desktop

clique

1

id_da_região

tipo de dispositivo

event_type

seq

1

iOS

clique

1

1

Android

conversão

1

2

iOS

clique

1

2

desktop

clique

1

Um fluxo REPLACE USING (region_id) SEQUENCE BY seq recebe essas atualizações para as regiões 1 e 3. A região 2 não tem atualizações:

id_da_região

tipo de dispositivo

event_type

seq

1

iOS

clique

2

1

Android

conversão

2

1

desktop

clique

2

3

iOS

clique

1

3

desktop

clique

2

id_da_região

tipo de dispositivo

event_type

seq

1

iOS

clique

2

1

Android

conversão

2

1

desktop

clique

2

3

iOS

clique

1

3

desktop

clique

2

O destino torna-se:

id_da_região

tipo de dispositivo

event_type

seq

Resultado

1

iOS

clique

2

Substituído, porque a seq 2 é maior que a seq 1

1

Android

conversão

2

Substituído, porque a seq 2 é maior que a seq 1

1

desktop

clique

2

Substituído, porque a seq 2 é maior que a seq 1

2

iOS

clique

1

Não alterado, porque a key não está presente nesta atualização

2

desktop

clique

1

Não alterado, porque a key não está presente nesta atualização

3

desktop

clique

2

Adicionado. A linha seq 1 para a região 3 não é adicionada, porque apenas a sequência mais alta para uma key é aplicada.

id_da_região

tipo de dispositivo

event_type

seq

Resultado

1

iOS

clique

2

Substituído, porque a seq 2 é maior que a seq 1

1

Android

conversão

2

Substituído, porque a seq 2 é maior que a seq 1

1

desktop

clique

2

Substituído, porque a seq 2 é maior que a seq 1

2

iOS

clique

1

Não alterado, porque a key não está presente nesta atualização

2

desktop

clique

1

Não alterado, porque a key não está presente nesta atualização

3

desktop

clique

2

Adicionado. A linha seq 1 para a região 3 não é adicionada, porque apenas a sequência mais alta para uma key é aplicada.

Requisitos

Os fluxos REPLACE USING têm os seguintes requisitos:

  • REPLACE USING flows executados no Databricks Runtime 18.2 e acima, em compute clássico ou serverless. A Databricks recomenda o Unity Catalog.
  • A fonte deve ser uma fonte de transmissão. REPLACE USING rejeita uma fonte que não seja de transmissão.
  • Você deve especificar pelo menos uma coluna de key e exatamente uma coluna SEQUENCE BY.

Quando usar fluxos REPLACE USING

Os LakeFlow Pipelines oferecem três fluxos que substituem linhas existentes. Escolha com base na aparência da sua origem e em como ela identifica as linhas a serem substituídas:

  • Use REPLACE USING quando sua origem for uma série de snapshots parciais indexados por coluna. O REPLACE USING substitui apenas os dados que possuem uma correspondência nos dados de entrada, deixando todos os outros dados inalterados. Não requer uma key primária.
  • Use o AUTO CDC quando sua origem for um feed de captura de dados de alterações (CDC) com operações explícitas de insert , update e delete , ou se você precisar de histórico de dimensão que muda lentamente (SCD) Tipo 2 . O AUTO CDC também requer uma chave primária verdadeira. Consulte As APIs do AUTO CDC: simplifique a captura de dados de alterações (CDC) com pipelines.
  • Use REPLACE WHERE quando sua origem for um snapshot e você quiser recomputar e substituir um intervalo da tabela de destino selecionado por um predicado, por exemplo, os últimos 7 dias, como uma operação em lote. Não requer uma chave primária. Consulte Processamento em lote com fluxos REPLACE WHERE.

Criar um fluxo REPLACE USING

Defina fluxos REPLACE USING em SQL ou Python.

Use a cláusula FLOW REPLACE USING em linha com CREATE STREAMING TABLE:

SQL
CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Alternativamente, use a sintaxe de formato longo CREATE FLOW:

SQL
CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
nota

BY NAME é necessário em SQL. Ele combina colunas por nome em vez de por posição.

Sequenciamento e dados fora de ordem

A coluna SEQUENCE BY torna o resultado independente da ordem em que as atualizações chegam. Uma linha é aplicada a uma chave apenas se sua sequência for maior do que a sequência já armazenada para essa chave; portanto, uma linha atrasada ou reproduzida que seja mais antiga do que o valor atual é ignorada. As chaves não presentes em uma atualização permanecem inalteradas.

Siga estas práticas para que a substituição se comporte de forma previsível:

Prática

Motivo

Use uma sequência que aumente estritamente por versão de key, como um Timestamp, número de versão ou deslocamento de Logs.

Duas linhas com a mesma key e a mesma sequência são mantidas, o que resulta em linhas duplicadas para essa key.

Use uma sequência não nula.

Uma sequência null pode levar a um comportamento indefinido.

Prática

Motivo

Use uma sequência que aumente estritamente por versão de key, como um Timestamp, número de versão ou deslocamento de Logs.

Duas linhas com a mesma key e a mesma sequência são mantidas, o que resulta em linhas duplicadas para essa key.

Use uma sequência não nula.

Uma sequência null pode levar a um comportamento indefinido.

Expectativas

Os fluxos REPLACE USING oferecem suporte a expectativas. warn e fail se comportam como em outros fluxos: warn mantém as linhas que violam a regra e registra a violação, e fail interrompe a atualização. Consulte Gerenciar a qualidade dos dados com expectativas de pipeline.

Uma expectativa drop trata uma linha violadora como se a origem nunca a tivesse produzido. A linha descartada não substitui, exclui ou modifica chaves correspondentes na tabela de destino:

  • O descarte ocorre antes da desduplicação, portanto, o fluxo mantém a versão válida mais recente para a key.
  • Se todas as linhas de entrada para uma key forem descartadas, as linhas existentes da key permanecerão inalteradas.
  • Como uma linha descartada não define um limite mínimo de sequência, uma atualização válida posterior ainda é processada, mesmo que sua sequência seja inferior à da linha descartada.

Limitações

Os fluxos REPLACE USING têm as seguintes limitações:

  • O REPLACE USING é compatível com um único fluxo por tabela de destino. Não é possível combinar o REPLACE USING com outro tipo de fluxo no mesmo destino.
  • A tabela de destino deve ser criada dentro do pipeline.
  • A fonte deve ser uma fonte de transmissão.
  • Você deve especificar pelo menos uma coluna de key e uma coluna SEQUENCE BY. As colunas de chave não podem ser repetidas, e o tipo de cada coluna de chave deve ser classificável. Tipos atômicos, como números inteiros, strings e datas, podem ser chaves, enquanto MAP e VARIANT não podem.

Exemplos

Os exemplos a seguir leem de samples.wanderbricks.booking_updates, uma tabela de amostra de alterações de estado de reserva que está disponível em todos os workspaces habilitados para o Unity Catalog. Cada reserva aparece uma vez por alteração, portanto booking_id se repete com um novo booking_update_id. Consulte dataset Wanderbricks.

Exemplo 1: Manter o registro mais recente para cada chave

Mantenha apenas o estado atual de cada reserva. O fluxo utiliza chaves em booking_id e sequencia por booking_update_id, portanto, a atualização mais recente de uma reserva substitui as anteriores. Use AUTO CDC em vez disso quando sua origem for um feed de alterações com operações explícitas de inserção, atualização e exclusão.

SQL
CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Este exemplo sequencia por booking_update_id em vez do timestamp updated_at, pois várias atualizações para a mesma reserva podem compartilhar um timestamp. As linhas que empatam na sequência são anexadas em vez de substituídas, o que deixaria mais de uma linha para essas reservas.

Exemplo 2: key em mais de uma coluna

Quando um registro é identificado por uma combinação de colunas, liste todas elas em REPLACE USING. Aqui, cada reserva é identificada por (property_id, booking_id), portanto, o fluxo mantém o estado atual de cada reserva por propriedade. Se uma coluna de key puder ser nula, REPLACE USING corresponde nulo a nulo em vez de ignorar a linha.

SQL
CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Exemplo 3: descartar registros inválidos com uma expectativa

Adicione uma expectativa para manter linhas incorretas fora do destino. Uma linha descartada é tratada como se a fonte nunca a tivesse produzido: ela não substitui nem exclui a key correspondente, e o fluxo retorna para a última linha válida para essa key. Este fluxo descarta atualizações que não possuem um total_amount positivo.

Python
from pyspark import pipelines as dp

@dp.table(
name="bookings_validated",
replace_using=["booking_id"],
sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
return spark.readStream.table("samples.wanderbricks.booking_updates")