replace_flow

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 USING tok a nelze ji kombinovat REPLACE USING s jiným typem toku, jako je připojený tok, automatický CDC tok nebo REPLACE WHERE tok.
  • Dotaz musí být streamovaným dotazem. @dp.replace_flow odmítá zdroj, který není streamován.
  • REPLACE USING toky vyžadují Databricks Runtime 18.2 a vyšší.