Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Important
Tato funkce je v beta verzi.
Dekorátor @dp.replace_flow vytváří tok NAHRAZENÍ POUŽITÍ pro streamovací tabulku ve vašem pipeline. Při každé aktualizaci flow nahrazuje všechny řádky v cílové tabulce, které odpovídají klíčovým sloupcům replace_using , a všechny ostatní řádky zůstávají nedotčené. Funkce musí vrátit streamingový datový rámec Apache Spark. Viz částečná výměna snímků pomocí NAHRAZENÍ POMOCÍ (REPLACEMENT USING flows).
Použijte @dp.replace_flow , když je zdroj série částečných snímků podle sloupců. Pro definování cílové tabulky a toku v jednom příkazu místo toho přepošlete replace_using a sequence_by do @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>)
Parametry
| Parameter | Typ | Description |
|---|---|---|
| funkce | function |
Required. Funkce, která vrací datový rámec streamování Apache Sparku z uživatelem definovaného dotazu. |
target |
str |
Required. Název streamovací tabulky, která je cílem toku. |
replace_using |
list |
Required. Klíčové sloupce, které určují, které cílové řádky nahradit. Uveďte alespoň jeden sloupec. Klíčové sloupce nelze opakovat a typ každého klíčového sloupce musí být tříditelný. |
sequence_by |
str nebo Column |
Required. Sloupec, který objednává aktualizace. Pro každý klíč vyhrává nejvyšší sekvence a řádek nižší sekvence nikdy nepřepíše vyšší posloupnost, která je již v cíli. |
name |
str |
Název toku. Pokud není zadaný, nastaví se výchozí hodnota názvu funkce. |
comment |
str |
Popis toku |
spark_conf |
dict |
Seznam konfigurací Sparku pro spuštění tohoto dotazu. |
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")
Používejte více než jeden klíčový sloupec, když je záznam identifikován kombinací sloupců:
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")
Omezení
- Streamovací tabulka podporuje jeden
REPLACE USINGtok a nelze ji kombinovatREPLACE USINGs jiným typem toku, jako je připojený tok, automatický CDC tok neboREPLACE WHEREtok. - Dotaz musí být streamovaným dotazem.
@dp.replace_flowodmítá zdroj, který není streamován. -
REPLACE USINGtoky vyžadují Databricks Runtime 18.2 a vyšší.