replace_flow

Important

Ez a funkció bétaverzióban érhető el.

A @dp.replace_flow dekorátor létrehoz egy HELYETTESÍTÉS HASZNÁLATOS folyamatot egy streaming táblához a csővezetékedben. Minden frissítésnél a folyamat lecseréli a céltáblában az összes olyan sort, amely megegyezik a replace_using kulcsoszlopokkal, és a többi sort érintetlenül hagyja. A függvénynek Apache Spark streaming DataFrame-et kell visszaadnia. Lásd Részleges snapshot helyettesítés CSERÉLJ USING flow-okkal.

Használd, @dp.replace_flow ha a forrásod részleges pillanatképek sorozata, amelyeket oszlopok szerint kulcsolnak. A céltáblát és az áramlást egyetlen utasításban definiálni helyette add replace_using át a @dp.táblánaksequence_by.

Szintaxis

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

Paraméter Típus Leírás
függvény function Required. Egy függvény, amely egy Apache Spark streamelési DataFrame-et ad vissza egy felhasználó által megadott lekérdezésből.
target str Required. Az áramlás célpontja a streaming tábla neve.
replace_using list Required. A kulcsoszlopok, amelyek meghatározzák, mely célsorokat kell helyettesíteni. Legalább egy oszlopot jelölj meg. A kulcsoszlopokat nem lehet ismételni, és minden kulcsoszlop típusának rendezhetőnek kell lennie.
sequence_by str vagy Column Required. Az oszlop, amely a frissítéseket rendeli. Minden kulcsnál a legmagasabb sorozat nyer, és egy alacsonyabb sorrendű sor soha nem írja felül a célban már meglévő magasabb sort.
name str A folyam neve. Ha nincs megadva, alapértelmezés szerint a függvény neve lesz.
comment str A folyamat leírása.
spark_conf dict A lekérdezés végrehajtásához szükséges Spark-konfigurációk listája.

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

Használj több kulcsoszlopot, ha egy rekordot oszlopok kombinációjával azonosítunk:

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

  • Egy streaming tábla egyetlen REPLACE USING flow-t támogat, és nem kombinálhat REPLACE USING más áramlástípussal, például append flow-val, automatikus CDC áramlással vagy áramlással REPLACE WHERE .
  • A lekérdezésnek streamelési lekérdezésnek kell lennie. @dp.replace_flow elutasítja a nem streaming forrást.
  • REPLACE USING a flow-khoz Databricks Runtime 18.2 vagy annál magasabb verziók szükségesek.