Substituição parcial de instantâneo com REPLACE USING fluxos

Importante

Este recurso está em versão Beta.

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

Uma coluna ordena as atualizações para que o resultado esteja correto mesmo quando as SEQUENCE BY atualizações chegam fora de ordem. Para cada chave, a sequência mais alta vence, e uma linha de sequência inferior nunca sobrescreve uma linha superior já no alvo. As linhas que partilham 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 event_type seq
1 iOS clicar 1
1 Android conversão 1
2 iOS clicar 1
2 ambiente de trabalho clicar 1

Um REPLACE USING (region_id) SEQUENCE BY seq fluxo recebe estas 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 clicar 2
1 Android conversão 2
1 ambiente de trabalho clicar 2
3 iOS clicar 1
3 ambiente de trabalho clicar 2

O alvo torna-se:

region_id device_type event_type seq Outcome
1 iOS clicar 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 ambiente de trabalho clicar 2 Substituído, porque o seq 2 é maior que o seq 1
2 iOS clicar 1 Inalterada, porque a chave não está presente nesta atualização
2 ambiente de trabalho clicar 1 Inalterada, porque a chave não está incluída nesta atualização
3 ambiente de trabalho clicar 2 Adicionado. A linha seq 1 para a região 3 não é adicionada, porque apenas a sequência mais alta para uma chave é aplicada.

Requisitos

Os fluxos de REPLACE USING têm os seguintes requisitos:

  • SUBSTITUIR UTILIZANDO fluxos executados no Databricks Runtime 18.2 e posteriores, em computação clássica ou serverless. O Databricks recomenda o Unity Catalog.
  • A fonte deve ser uma fonte em streaming. REPLACE USING rejeita uma origem não contínua.
  • Deve especificar pelo menos uma coluna-chave e exatamente uma SEQUENCE BY coluna.

Quando usar REPLACE USING em fluxos

Os oleodutos de fluxo de lago oferecem três fluxos que sobrescrevem as linhas existentes. Escolha com base no aspeto da sua fonte e em como identifica as linhas a substituir:

  • Use REPLACE USING quando a sua origem for uma série de instantâneos parciais identificados por coluna. SUBSTITUIR UTILIZANDO sobrescreve apenas os dados que têm correspondência nos dados de entrada, deixando todos os outros dados inalterados. Não requer uma chave primária.
  • Utilize AUTO CDC quando a sua origem for um fluxo de captura de dados alterados (CDC) com operações explícitas de inserção, atualização e eliminação, ou quando precisar de histórico de dimensão de evolução lenta (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ção com pipelines.
  • Utilize REPLACE WHERE quando a sua origem for um instantâneo e quiser recalcular e substituir um intervalo da tabela de destino selecionado através de um predicado, por exemplo, os últimos 7 dias, numa operação em lote. Não requer uma chave primária. Consulte o processamento em lote com fluxos REPLACE WHERE.

Crie um fluxo de substituir utilizando

Defina fluxos REPLACE USING em SQL ou em Python.

SQL

Utilize a cláusula FLOW REPLACE USING na mesma linha que 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);

Alternativamente, use a sintaxe de formato longo CREATE FLOW :

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 é exigido em SQL. Faz a correspondência das colunas com base no nome e não pela posição.

Python

Declare a tabela e o fluxo juntamente 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")

Alternativamente, especifique 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 está definida.

Sequenciação 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 apenas se a sua sequência for maior do que a sequência já armazenada para essa chave, pelo que uma linha tardia ou reproduzida e mais antiga do que o valor atual é ignorada. As teclas que não estão presentes numa atualização ficam intocadas.

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

Practice Justificação
Utilize uma sequência que aumente de forma estrita para cada versão da chave, como uma marca temporal, um número de versão ou um offset do registo. Duas linhas com a mesma chave e a mesma sequência são ambas mantidas, o que resulta em linhas duplicadas para essa chave.
Usa uma sequência não nula. Uma sequência nula pode levar a comportamentos indefinidos.

Expectations

Substituir com fluxos cumpre as expectativas. warn e fail comportam-se como nos outros fluxos: warn continuam a violar as linhas e registam a violação, e fail interrompem a atualização. Consulte Gerir a qualidade dos dados com as expectativas do fluxo de dados.

Uma drop expectativa considera uma linha em violação como se a origem nunca a tivesse produzido. A linha eliminada não substitui, elimina ou modifica as chaves correspondentes na tabela de destino:

  • O descarte acontece antes da desduplicação, por isso o fluxo mantém a versão válida mais recente para a chave.
  • Se todas as linhas de entrada de uma chave forem eliminadas, as linhas existentes da chave ficam intocadas.
  • Como uma linha caída não define um piso de sequência, uma atualização válida posterior ainda assim aparece mesmo que a sua sequência seja inferior à da linha descartada.

Limitações

Os fluxos «Substituir utilizando» têm as seguintes limitações:

  • REPLACE USING suporta um único fluxo por tabela de destino. A combinação de REPLACE USING com outro tipo de fluxo no mesmo destino não é suportada.
  • A tabela alvo deve ser criada dentro do pipeline.
  • A fonte deve ser uma fonte em streaming.
  • Deve especificar pelo menos uma coluna-chave e uma coluna SEQUENCE BY. As colunas chave não podem ser repetidas, e o tipo de cada coluna chave deve ser ordenável. Tipos atómicos, como inteiros, cadeias de carateres e datas, podem ser chaves, ao passo que MAP e VARIANT não podem.
  • Para tabelas de streaming autónomas, consulte Aplicar a substituição parcial de instantâneo com fluxos REPLACE USING para conhecer as diferenças de sintaxe.

Examples

Os exemplos seguintes leem de samples.wanderbricks.booking_updates, uma tabela de exemplo de alterações de estado de reservas disponível em todos os espaços de trabalho habilitados pelo Unity Catalog. Cada reserva aparece uma vez por cada alteração, pelo que booking_id se repete com um novo booking_update_id. Consulte o conjunto de dados Wanderbricks.

Exemplo 1: Manter o registo mais recente de cada tecla

Mantém apenas o estado atual de cada reserva. O fluxo é indexado por booking_id e sequenciado por booking_update_id, pelo que a atualização mais recente de uma reserva substitui as anteriores. Use o AUTO CDC em vez disso quando a sua fonte for um feed de alterações com operações explícitas de inserção, atualização e eliminaçã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 é ordenado por booking_update_id em vez do carimbo temporal updated_at, porque várias atualizações da mesma reserva podem partilhar o mesmo carimbo temporal. As linhas que têm a mesma sequência são acrescentadas em vez de serem substituídas, o que resulta em mais de uma linha para essas reservas.

Exemplo 2: Chave em mais do que uma coluna

Quando um registo é identificado por uma combinação de colunas, liste-os todos em REPLACE USING. Aqui cada reserva é identificada por (property_id, booking_id), pelo que o fluxo mantém o estado atual de cada reserva por propriedade. Se uma coluna-chave puder ter valor nulo, REPLACE USING faz corresponder valores nulos entre si, 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: Eliminar registos inválidos com uma expectativa

Adicione uma regra de validação para impedir que linhas inválidas cheguem ao destino. Uma linha caída é tratada como se a fonte nunca a tivesse produzido: não substitui nem elimina a chave correspondente, e o fluxo volta para a linha válida mais recente dessa chave. Este fluxo descarta atualizações que não têm 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")