Pular para o conteúdo principal

Substituição parcial de Snapshot com fluxos REPLACE USING

info

Beta

Esse recurso está em Beta.

REPLACE USING transmissões em LakeFlow Pipelines mantêm uma tabela atualizada a partir de uma transmissão de snapshots parciais. Elas substituem todas as linhas que correspondem às colunas de key especificadas e deixam o resto da tabela inalterado. Elas se ajustam a fontes que enviam periodicamente um conjunto completo de linhas para uma key, como arquivos recarregados quando seus conteúdos são alterados.

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. As colunas de key não podem ser repetidas e seus tipos devem ser ordenáveis. Tipos atômicos, como números inteiros, strings e datas, podem ser keys. MAP e VARIANT não podem ser keys.

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.

nota

Para tabelas de transmissão autônomas, consulte Apply partial Snapshot replacement with REPLACE USING flows para obter diferenças de sintaxe.

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.

Operações suportadas​

As seguintes operações são suportadas na query do fluxo:

  • Projeção de coluna: selecione ou reordene um subconjunto de colunas.
  • Filtros com WHERE.
  • Expressões escalares, como CAST, aritméticas e CASE.
  • Deduplicação com SELECT DISTINCT.
  • Limites de linhas com LIMIT.
  • Funções geradoras, como EXPLODE e POSEXPLODE.
  • Uniões de duas fontes de transmissão com UNION ALL.
  • Joins estáticos de transmissão (stream-static joins): internas e externas à esquerda, com a transmissão no lado esquerdo.
  • Transmissão-transmissão joins internas.

As seguintes operações exigem configuração adicional:

  • Janelas baseadas em tempo, como window(ts, '5 minutes'), exigem uma marca d'água.
  • As junções externas de transmissão-transmissão (stream-stream outer joins) exigem uma marca d'água e uma condição de intervalo de tempo.

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.

As seguintes operações não são aceitas na query do fluxo:

  • Agregações, como SUM, COUNT e GROUP BY.
  • Funções de janela que não sejam janelas baseadas em tempo, como ROW_NUMBER() OVER (...), mesmo com uma marca d'água.
  • Ordenando com ORDER BY.
  • Operações de conjunto (set operations), como INTERSECT e EXCEPT.
  • Uniões que misturam uma fonte de transmissão com uma fonte sem transmissão.
  • Leitura de uma fonte que não é de transmissão, como spark.range().
  • Junções externas de transmissão-static full e right.

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