Substituição parcial de instantâneos por fluxos REPLACE USING

Importante

Esse recurso está em Beta.

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

Uma SEQUENCE BY coluna ordena as atualizações para que o resultado esteja correto mesmo quando as atualizações chegam fora de ordem. Para cada chave, a sequência mais alta prevalece, e uma linha de menor sequência jamais substitui uma linha com sequência mais alta já presente no destino. Linhas que compartilham a mesma chave e a mesma sequência são adicionadas em vez de substituídas.

Como funciona REPLACE USING

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

region_id device_type tipo_de_evento seq
1 iOS clique 1
1 Android conversão 1
2 iOS clique 1
2 área de trabalho clique 1

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

region_id device_type event_type seq
1 iOS clique 2
1 Android conversão 2
1 área de trabalho clique 2
3 iOS clique 1
3 área de trabalho clique 2

O alvo se torna:

region_id device_type tipo_de_evento seq Resultado
1 iOS clique 2 Substituído, porque o seq 2 é maior que o seq 1
1 Android conversão 2 Substituído, porque o seq 2 é maior que o seq 1
1 área de trabalho clique 2 Substituído, porque o seq 2 é maior que o seq 1
2 iOS clique 1 Inalterado, porque a chave não está presente nesta atualização
2 área de trabalho clique 1 Inalterado, porque a chave não está presente nesta atualização
3 área de trabalho clique 2 Adicionado. A linha de sequência 1 da região 3 não é adicionada, porque apenas a sequência mais alta de uma chave é aplicada.

Requisitos

Os fluxos REPLACE USING têm os seguintes requisitos:

  • Os fluxos REPLACE USING são executados no Databricks Runtime 18.2 ou superior, em computação sem servidor ou clássica. O Databricks recomenda o Unity Catalog.
  • A origem deve ser uma fonte de streaming. REPLACE USING rejeita uma fonte não de streaming.
  • Você deve especificar pelo menos uma coluna-chave e exatamente uma SEQUENCE BY coluna.

Quando usar fluxos REPLACE USING

Os pipelines do Lakeflow oferecem três fluxos que substituem linhas existentes. Escolha com base em como é a sua fonte e como ela identifica as linhas a substituir:

  • Use REPLACE USING quando sua origem for uma série de instantâneos parciais codificados por coluna. REPLACE USING sobrescreve apenas os dados que encontram correspondência nos dados de entrada, deixando todos os outros dados inalterados. Não precisa de uma chave primária.
  • Use AUTO CDC quando a fonte for um feed da captura de dados de alterações (CDC) com operações explícitas de inserção, atualização e exclusão, ou quando você precisar do histórico de dimensões de alterações lentas (SCD) tipo 2. O AUTO CDC também exige uma chave primária verdadeira. Consulte As APIs AUTO CDC: Simplifique a captura de dados de alterações com pipelines.
  • Use REPLACE WHERE quando a origem for um instantâneo e você quiser recalcular e substituir um intervalo da tabela de destino selecionado por um predicado, por exemplo, os últimos 7 dias, em uma operação em lote. Não precisa de uma chave primária. Consulte Processamento em lote com fluxos REPLACE WHERE.

Criar um fluxo REPLACE USING

Defina fluxos de trabalho REPLACE USING em SQL ou Python.

SQL

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

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);

Como alternativa, use a sintaxe de forma CREATE FLOW longa:

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);

Note

BY NAME é necessário no SQL. Faz a correspondência das colunas pelo nome em vez da posição.

Python

Declare a tabela e o fluxo junto com @dp.table:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

Como alternativa, direcione para uma tabela de streaming existente com @dp.replace_flow:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using é uma lista de colunas-chave. sequence_by é um nome de coluna ou uma Column expressão, e é necessária sempre que replace_using é definida.

Sequenciamento e dados fora de ordem

A SEQUENCE BY coluna torna o resultado independente da ordem em que as atualizações chegam. Uma linha é aplicada a uma chave somente se sua sequência for maior do que a sequência já armazenada para essa chave, então uma linha tardia ou reproduzida que seja mais antiga que o valor atual é ignorada. Teclas que não estão presentes em uma atualização ficam intocadas.

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 da chave, como carimbo de data/hora, número da versão ou deslocamento do log. Duas linhas com a mesma chave e a mesma sequência são mantidas, o que resulta em linhas duplicadas para essa chave.
Use uma sequência não nula. Uma sequência nula pode levar a comportamentos indefinidos.

Expectations

Os fluxos REPLACE USING atendem a expectativas. warn e fail se comportam como em outros fluxos: warn continuam violando as linhas, registram a violação e fail interrompem a atualização. Confira Gerenciar a qualidade dos dados com as expectativas do pipeline.

Uma expectativa drop trata uma linha em violação como se a fonte jamais a tivesse produzido. A linha eliminada não substitui, exclui ou modifica as chaves correspondentes na tabela de destino:

  • O descarte acontece antes da eliminação de duplicação, logo, o fluxo mantém a versão válida mais recente para a chave.
  • Se todas as linhas de entrada de uma chave forem descartadas, as linhas existentes dessa chave permanecerão inalteradas.
  • Como uma linha descartada não estabelece um limite mínimo de sequência, uma atualização válida posterior ainda é aplicada, mesmo que tenha uma sequência inferior à da linha descartada.

Limitações

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

  • REPLACE USING suporta um único fluxo por tabela de destino. Não há suporte para combinar REPLACE USING com outro tipo de fluxo no mesmo destino.
  • A tabela de destino deve ser criada no pipeline.
  • A origem deve ser uma fonte de streaming.
  • Você deve especificar pelo menos uma coluna-chave e uma coluna SEQUENCE BY. Colunas chave não podem ser repetidas, e o tipo de cada coluna chave deve ser ordenável. Tipos atômicos, como inteiros, cadeias e datas, podem ser chaves, enquanto MAP e VARIANT não podem.
  • Para tabelas de streaming autônomo, consulte Aplicar substituição de instantâneo parcial por fluxos REPLACE USING para obter as diferenças de sintaxe.

Exemplos

Os exemplos a seguir leem dados de samples.wanderbricks.booking_updates, uma tabela de amostra de alterações no estado das reservas disponível em todos o workspace habilitado para o Catálogo do Unity. Cada reserva aparece uma vez a cada alteração, logo, booking_id se repete a cada novo booking_update_id. Veja o conjunto de dados Wanderbricks.

Exemplo 1: Manter o registro mais recente de cada chave

Mantenha apenas o estado atual de cada reserva. O fluxo usa booking_id como chave e sequencia por booking_update_id, de maneira que a atualização mais recente de uma reserva substitui as anteriores. Use o AUTO CDC quando sua fonte for um feed de mudanças 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);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Este exemplo ordena por booking_update_id, em vez do carimbo de data/hora updated_at, porque várias atualizações de uma mesma reserva podem compartilhar o mesmo carimbo de data/hora. As linhas na mesma sequência são acrescentadas, em vez de substituídas, o que faz com que haja mais de uma linha para esses agendamentos.

Exemplo 2: Chave em mais de uma coluna

Quando um registro é identificado por uma combinação de colunas, liste-as todas em REPLACE USING. Aqui, cada reserva é identificada por (property_id, booking_id), então o fluxo mantém o estado atual de cada reserva por propriedade. Se uma coluna-chave puder ser nula, REPLACE USING combinará nulo com 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);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Exemplo 3: Descartar registros inválidos com uma expectativa

Adicione uma validação para impedir que linhas inválidas cheguem ao destino. Uma linha caída é tratada como se a fonte nunca a tivesse produzido: ela não substitui nem exclui a chave correspondente, e o fluxo volta para a linha válida mais recente dessa chave. Esse fluxo descarta atualizações que não têm um total_amount positivo.

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")