Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
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 BYcoluna.
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 queMAPeVARIANTnã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")