Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
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 USINGflow-t támogat, és nem kombinálhatREPLACE USINGmás áramlástípussal, például append flow-val, automatikus CDC áramlással vagy áramlássalREPLACE WHERE. - A lekérdezésnek streamelési lekérdezésnek kell lennie.
@dp.replace_flowelutasítja a nem streaming forrást. -
REPLACE USINGa flow-khoz Databricks Runtime 18.2 vagy annál magasabb verziók szükségesek.