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.
Il @dp.replace_flow decoratore crea un flusso SOSTITUISCE USANDO per un tavolo di streaming nella tua pipeline. Ad ogni aggiornamento, il flusso sostituisce tutte le righe della tabella di destinazione che corrispondono alle replace_using colonne chiave e lascia tutte le altre righe intatte. La funzione deve restituire un dataframe di streaming Apache Spark.
Vedi Sostituzione parziale di snapshot con SOSTITUIRE i flussi USING (SOSTITUIRE UTILIZZANDO i flows).
Usalo @dp.replace_flow quando la tua sorgente è una serie di istantanee parziali chiaviate per colonna. Per definire la tabella di destinazione e il flusso in un'unica istruzione, passa replace_using e sequence_byva a @dp.tabella.
Sintassi
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.replace_flow(
target = "<target-table-name>",
replace_using = ["<key-column>", "<key-column>"],
sequence_by = "<sequence-column>",
name = "<flow-name>", # optional, defaults to function name
comment = "<comment>", # optional
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}) # optional
def <function-name>():
return (<streaming-query>)
Parameters
| Parametro | Tipo | Description |
|---|---|---|
| funzione | function |
Required. Funzione che restituisce un dataframe di streaming Apache Spark da una query definita dall'utente. |
target |
str |
Required. Il nome della tabella di streaming che è il bersaglio del flow. |
replace_using |
list |
Required. Le colonne chiave che identificano quali righe target sostituire. Specifica almeno una colonna. Le colonne chiave non possono essere ripetute e il tipo di ogni colonna chiave deve essere sorrettibile. |
sequence_by |
str oppure Column |
Required. La colonna che ordina gli aggiornamenti. Per ogni chiave, vince la sequenza più alta, e una riga di sequenza inferiore non sovrascrive mai una più alta già nel target. |
name |
str |
Nome del flusso. Se non specificato, per impostazione predefinita viene impostato il nome della funzione. |
comment |
str |
Descrizione del flusso. |
spark_conf |
dict |
Elenco delle configurazioni di Spark per l'esecuzione di questa query. |
Examples
from pyspark import pipelines as dp
# Keep the latest row for each order from a stream of partial snapshots
dp.create_streaming_table("orders_current")
@dp.replace_flow(
target = "orders_current",
replace_using = ["order_id"],
sequence_by = "updated_at"
)
def orders_flow():
return spark.readStream.table("order_updates")
Usa più di una colonna chiave quando un record è identificato da una combinazione di colonne:
from pyspark import pipelines as dp
dp.create_streaming_table("accounts_current")
@dp.replace_flow(
target = "accounts_current",
replace_using = ["region", "account_id"],
sequence_by = "updated_at"
)
def accounts_flow():
return spark.readStream.table("account_updates")
Limitations
- Una tabella di flusso supporta un singolo
REPLACE USINGflusso e non può combinarsiREPLACE USINGcon un altro tipo di flusso come un flusso append, un flusso automatico CDC o unREPLACE WHEREflusso di corrente. - La query deve essere una query di streaming.
@dp.replace_flowrifiuta una fonte non in streaming. -
REPLACE USINGi flussi richiedono Databrick Runtime 18.2 e superiori.