Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Importante
Questa funzionalità è in versione beta.
Un flusso REPLACE USING mantiene una tabella di destinazione sincronizzata con un'origine dati in streaming: sostituisce tutte le righe che corrispondono alle colonne chiave specificate e lascia invariati tutti gli altri dati.
Una SEQUENCE BY colonna ordina gli aggiornamenti in modo che il risultato sia corretto anche quando gli aggiornamenti arrivano fuori ordine. Per ogni chiave, prevale la sequenza più alta e una riga con sequenza inferiore non sovrascrive mai quella con sequenza più alta già presente nella destinazione. Le righe che condividono la stessa chiave e la stessa sequenza vengono aggiunte invece che sostituite.
Come funziona REPLACE USING
Consideriamo una tabella degli eventi che contiene eventi di clic e conversione per due regioni, sequenziati da:seq
| region_id | device_type | tipo_evento | seq |
|---|---|---|---|
| 1 | iOS | click | 1 |
| 1 | Android | conversione | 1 |
| 2 | iOS | click | 1 |
| 2 | Desktop | click | 1 |
Un REPLACE USING (region_id) SEQUENCE BY seq flusso riceve questi aggiornamenti per le regioni 1 e 3. La Regione 2 non ha aggiornamenti:
| region_id | device_type | event_type | seq |
|---|---|---|---|
| 1 | iOS | click | 2 |
| 1 | Android | conversione | 2 |
| 1 | Desktop | click | 2 |
| 3 | iOS | click | 1 |
| 3 | Desktop | click | 2 |
L'obiettivo diventa:
| region_id | device_type | event_type | seq | Risultato |
|---|---|---|---|---|
| 1 | iOS | click | 2 | Sostituito, perché il seq 2 è maggiore del seq 1 |
| 1 | Android | conversione | 2 | Sostituito, perché il seq 2 è maggiore del seq 1 |
| 1 | Desktop | click | 2 | Sostituito, perché il seq 2 è maggiore del seq 1 |
| 2 | iOS | click | 1 | Invariato, perché la chiave non è presente in questo aggiornamento |
| 2 | Desktop | click | 1 | Invariato, perché la chiave non è presente in questo aggiornamento |
| 3 | Desktop | click | 2 | Aggiunto. La riga con sequenza 1 per la regione 3 non viene aggiunta, perché per una chiave viene applicata solo la sequenza più alta. |
Requirements
I flussi di lavoro REPLACE USING hanno i seguenti requisiti:
- SOSTITUIRE USANDO i flussi eseguiti su Databrick runtime 18.2 e superiori, su calcolo classico o serverless. Databricks consiglia Unity Catalog.
- L'origine deve essere un'origine di streaming. REPLACE USING rifiuta una sorgente non di streaming.
- Devi specificare almeno una colonna chiave e esattamente una
SEQUENCE BYcolonna.
Quando usare REPLACE USING con i flussi
Le pipeline Lakeflow mettono a disposizione tre flussi che sovrascrivono le righe esistenti. Scegli in base a come appare la tua fonte e a come identifica le righe da sostituire:
- Usa REPLACE USING quando la sorgente è una serie di snapshot parziali indicizzati per colonna. REPLACE USING sovrascrive solo i dati che hanno una corrispondenza con i dati in arrivo, lasciando tutti gli altri dati intatti. Non richiede una chiave primaria.
- Usa AUTO CDC quando la tua sorgente è un feed di cattura dati di modifica (CDC) con operazioni esplicite di inserimento, aggiornamento e cancellazione , oppure quando hai bisogno di una cronologia di dimensione a cambiamento lento (SCD) di tipo 2 . AUTO CDC richiede anche una vera chiave primaria. Consulta le API AUTO CDC: semplificare la cattura dei dati modificati con le pipeline.
- Usa REPLACE WHERE quando la tua sorgente è una snapshot e vuoi ricalcolare e sovrascrivere un intervallo della tabella target selezionata da un predicato, ad esempio gli ultimi 7 giorni, come operazione batch. Non richiede una chiave primaria. Vedi Elaborazione batch con flussi REPLACE WHERE.
Crea un flusso SOSTITUIRE USANDO
Definisci i flow REPLACE USING in SQL o in Python.
SQL
Usa la clausola FLOW REPLACE USING in linea con 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);
In alternativa, usare la sintassi long-form 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 è obbligatorio in SQL. Corrisponde alle colonne in base al nome anziché alla posizione.
Python
Dichiara la tabella e il flusso insieme con @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")
In alternativa, scegli come destinazione una tabella di streaming esistente con @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 è un elenco di colonne chiave.
sequence_by è un nome di colonna o un'espressione Column , ed è richiesta ogni volta che replace_using è impostato.
Sequenziamento e dati non ordinati
La SEQUENCE BY colonna rende il risultato indipendente dall'ordine in cui arrivano gli aggiornamenti. Una riga viene applicata a una chiave solo se la sua sequenza è maggiore di quella già memorizzata per quella chiave, quindi viene ignorata una riga in ritardo o ripetuta e più vecchia del valore corrente. Le chiavi non presenti in un aggiornamento restano intatte.
Segui queste pratiche affinché la sostituzione si comporti in modo prevedibile:
| Pratica | Ragione |
|---|---|
| Usa una sequenza che aumenti in modo rigorosamente crescente per ciascuna versione della chiave, come un timestamp, un numero di versione o un offset del log. | Due righe con la stessa chiave e la stessa sequenza vengono entrambe mantenute, il che comporta righe duplicate per quella chiave. |
| Usa una sequenza non nulla. | Una sequenza nulla può portare a comportamenti indefiniti. |
Expectations
SOSTITUIRE UTILIZZANDO i flussi supportano le aspettative.
warn e fail si comportano come fanno sugli altri flow: warn continua a violare le righe, registra la violazione e fail interrompe l'aggiornamento. Vedi Gestisci la qualità dei dati con le aspettative della pipeline.
Una drop aspettativa tratta una riga in violazione come se l'origine non l'avesse mai prodotta. La riga eliminata non sostituisce, elimina o modifica le chiavi corrispondenti nella tabella di destinazione:
- Il drop avviene prima della deduplicazione, quindi il flow mantiene l'ultima versione valida della chiave.
- Se ogni riga in ingresso di una chiave viene eliminata, le righe esistenti della chiave restano intatte.
- Poiché una riga eliminata non stabilisce un piano di sequenza, un aggiornamento valido successivo viene comunque valido anche se la sua sequenza è inferiore a quella della riga eliminata.
Limitations
SOSTITUIRE UTILIZZANDO i flussi presenta le seguenti limitazioni:
- REPLACE USING supporta un singolo flusso per tabella di destinazione. La combinazione di REPLACE USING con un altro tipo di flusso nella stessa destinazione non è supportata.
- La tabella di destinazione deve essere creata all'interno della pipeline.
- L'origine deve essere un'origine di streaming.
- Devi specificare almeno una colonna chiave e una
SEQUENCE BYcolonna. Le colonne chiave non possono essere ripetute e il tipo di ogni colonna chiave deve essere ordinabile. I tipi atomici, come interi, stringhe e date, possono essere chiavi, mentreMAPeVARIANTnon possono. - Per le tabelle di streaming autonome, vedere Applicare la sostituzione parziale dello snapshot con i flussi REPLACE USING per le differenze di sintassi.
Examples
I seguenti esempi sono letti da samples.wanderbricks.booking_updates, una tabella di esempio delle modifiche allo stato delle prenotazioni disponibile in ogni spazio di lavoro abilitato al Catalogo Unity. Ogni prenotazione appare una volta per ogni modifica, quindi booking_id si ripete con un nuovo booking_update_id. Vedi il dataset di Wanderbricks.
Esempio 1: Conserva l'ultimo record per ogni chiave
Conserva solo lo stato attuale di ogni prenotazione. Il flow usa come chiave booking_id e ordina in sequenza in base a booking_update_id, quindi l'aggiornamento più recente di una prenotazione sostituisce i precedenti aggiornamenti. Usa invece AUTO CDC quando la tua sorgente è un feed di modifiche con operazioni esplicite di inserimento, aggiornamento e cancellazione.
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")
Questo esempio ordina in base a booking_update_id anziché al timestamp updated_at, perché diversi aggiornamenti della stessa prenotazione possono condividere un timestamp. Le righe che si legano alla sequenza vengono aggiunte invece che sostituite, il che lascerebbe più di una riga per quelle prenotazioni.
Esempio 2: Chiave su più di una colonna
Quando un record viene identificato da una combinazione di colonne, elencale tutte in REPLACE USING. Qui ogni prenotazione è identificata da (property_id, booking_id), così il flusso mantiene lo stato attuale di ogni prenotazione per proprietà. Se una colonna chiave può essere NULL, REPLACE USING fa corrispondere NULL a NULL invece di ignorare la riga.
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")
Esempio 3: Eliminare i record non validi con un'aspettativa
Aggiungi l'aspettativa di tenere fuori dal bersaglio le cattive linee. Una riga eliminata viene trattata come se la sorgente non l'avesse mai prodotta: non sostituisce né elimina la chiave corrispondente, e il flusso torna alla riga valida più recente per quella chiave. Questo flusso ignora gli aggiornamenti che non hanno un 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")