Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
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 USINGFlow und kann sich nicht mit einem anderen Flow-Typ wie einem Append Flow, einem Auto-CDC-Flow oder einemREPLACE WHEREFlow kombinierenREPLACE USING. - Die Abfrage muss eine Streamingabfrage sein.
@dp.replace_flowlehnt eine nicht-streamende Quelle ab. -
REPLACE USINGflows erfordern Databricks Runtime 18.2 und höher.