replace_flow

Important

Dieses Feature befindet sich in der Betaversion.

Der Dekorateur @dp.replace_flow erstellt einen ERSETZEN-NUTZUNGS-Flow für eine Streaming-Tabelle in deiner Pipeline. Bei jedem Update ersetzt der Flow alle Zeilen in der Zieltabelle, die mit den replace_using Schlüsselspalten übereinstimmen, und lässt alle anderen Zeilen unberührt. Die Funktion muss einen Apache Spark Streaming DataFrame zurückgeben. Siehe Teilweise Snapshot-Ersatz mit ERSETZEN UNTER Verwendung von Flows.

Verwenden Sie @dp.replace_flow , wenn Ihre Quelle eine Reihe von teilweisen Schnappschüssen ist, die nach Spalten verschlüsselt sind. Um die Zieltabelle und den Fluss in einer einzigen Anweisung zu definieren, übergebe replace_using stattdessen und sequence_by an @dp.table.

Syntax

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

Parameter

Parameter Typ Description
Funktion function Required. Eine Funktion, die einen Apache Spark Streaming DataFrame aus einer benutzerdefinierten Abfrage zurückgibt.
target str Required. Der Name der Streaming-Tabelle, die das Ziel des Flows ist.
replace_using list Required. Die Schlüsselspalten, die angeben, welche Zielzeilen ersetzt werden sollen. Geben Sie mindestens eine Spalte an. Schlüsselspalten können nicht wiederholt werden, und der Typ jeder Schlüsselspalte muss sortierbar sein.
sequence_by str oder Column Required. Die Spalte, die die Aktualisierungen anordnet. Für jeden Schlüssel gewinnt die höchste Sequenz, und eine niedrigere Zeile überschreibt niemals eine höhere, die bereits im Ziel ist.
name str Der Flussname. Wenn nicht angegeben, wird standardmäßig der Funktionsname verwendet.
comment str Eine Beschreibung für den Ablauf.
spark_conf dict Eine Liste der Spark-Konfigurationen für die Ausführung dieser Abfrage.

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

Verwenden Sie mehr als eine Schlüsselspalte, wenn ein Datensatz durch eine Kombination von Spalten identifiziert wird:

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

Einschränkungen

  • Eine Streaming-Tabelle unterstützt einen einzelnen REPLACE USING Flow und kann sich nicht mit einem anderen Flow-Typ wie einem Append Flow, einem Auto-CDC-Flow oder einem REPLACE WHERE Flow kombinierenREPLACE USING.
  • Die Abfrage muss eine Streamingabfrage sein. @dp.replace_flow lehnt eine nicht-streamende Quelle ab.
  • REPLACE USING flows erfordern Databricks Runtime 18.2 und höher.