replace_flow

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 USING flusso e non può combinarsi REPLACE USING con un altro tipo di flusso come un flusso append, un flusso automatico CDC o un REPLACE WHERE flusso di corrente.
  • La query deve essere una query di streaming. @dp.replace_flow rifiuta una fonte non in streaming.
  • REPLACE USING i flussi richiedono Databrick Runtime 18.2 e superiori.