Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
Important
Deze functie bevindt zich in de bètaversie.
De @dp.replace_flow decorateur maakt een VERVANG GEBRUIK-flow voor een streamingtabel in je pipeline. Bij elke update vervangt de flow alle rijen in de doeltabel die overeenkomen met de replace_using sleutelkolommen en laat alle andere rijen onaangeroerd. De functie moet een Apache Spark-streaming dataframe retourneren. Zie Gedeeltelijke snapshotvervanging met VERVANG MET stromen.
Gebruik @dp.replace_flow wanneer je bron een reeks gedeeltelijke snapshots is die per kolom zijn gecodeerd. Om de doel-tabel en de flow in één enkele instructie te definiëren, geef replace_using en sequence_by dan door naar @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>)
Parameters
| Parameter | Type | Description |
|---|---|---|
| function | function |
Required. Een functie die een Streaming DataFrame van Apache Spark retourneert vanuit een door de gebruiker gedefinieerde query. |
target |
str |
Required. De naam van de streamingtabel die het doel van de flow is. |
replace_using |
list |
Required. De sleutelkolommen die aangeven welke doelrijen vervangen moeten worden. Specificeer ten minste één kolom. Sleutelkolommen kunnen niet worden herhaald, en het type van elke sleutelkolom moet sorteerbaar zijn. |
sequence_by |
str of Column |
Required. De kolom die de updates bestelt. Voor elke sleutel wint de hoogste reeks, en een lagere reeks rij overschrijft nooit een hogere die al in het doel staat. |
name |
str |
De naam van de stroom. Als deze niet is opgegeven, wordt standaard de functienaam gebruikt. |
comment |
str |
Een beschrijving voor het proces. |
spark_conf |
dict |
Een lijst met Spark-configuraties voor de uitvoering van deze query. |
Voorbeelden
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")
Gebruik meer dan één sleutelkolom wanneer een record wordt geïdentificeerd door een combinatie van kolommen:
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
- Een streamingtabel ondersteunt één enkele
REPLACE USINGflow en kan niet gecombineerdREPLACE USINGworden met een ander flowtype zoals een appendflow, een automatische CDC-flow of eenREPLACE WHEREflow. - De query moet een streamingquery zijn.
@dp.replace_flowWijst een niet-streaming bron af. -
REPLACE USINGflows vereisen Databricks Runtime 18.2 en hoger.