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.
A @dp.append_flow dekorátor hozzáfűző folyamatokat vagy visszatöltéseket hoz létre a folyamattáblákhoz. A függvénynek Apache Spark streaming DataFrame-et kell visszaadnia. Lásd az adatok növekményes betöltését és feldolgozását a Lakeflow-folyamatokkal.
A hozzáfűzési folyamok megcélozhatnak streamelési táblákat vagy adatfogadókat.
Szemantika
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.append_flow(
target = "<target-table-name>",
name = "<flow-name>", # optional, defaults to function name
once = False, # optional
spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
comment = "<comment>", # optional
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>) #
Paraméterek
| Paraméter | Típus | Description |
|---|---|---|
| függvény | function |
Szükséges. 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 |
Szükséges. A hozzáfűzési folyamat céltáblájának vagy kimenetének neve. |
name |
str |
Az adatfolyam neve. Ha nincs megadva, alapértelmezés szerint a függvény neve lesz. |
once |
bool |
Igény szerint definiálja a folyamatot egyszeri folyamatként, például visszatöltésként. A once=True használata kétféleképpen változtatja meg a folyamatot:
|
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 |
import_checkpoint |
str |
Az út egy meglévő Strukturált Streaming ellenőrzőponthoz importálható be az áramlásba, így a migrált stream folytatja az utolsó elkötelezett eltolását, ahelyett, hogy újradolgozná a forrást. Az ellenőrzőpont importálása a Bétában van. Lásd : Strukturált streaming ellenőrzőpont migrálása. |
Példák
from pyspark import pipelines as dp
# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})
# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
return <streaming-query>
# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
return (
spark.read
.format("json")
.load("/path/to/backfill/")
)
# Create a Kafka sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "host:port",
"topic": "my_topic"
}
)
# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))
Migrálj egy strukturált streaming ellenőrzőpontot
Important
Az ellenőrzőpont importálása a Bétában van.
Használd import_checkpoint egy meglévő strukturált streaming munkaterhelés átmigrálásához egy csővezetékre anélkül, hogy újra feldolgozná a forrást. Állítsd be a checkpointLocation Strukturált Streaming lekérdezésedre, ami lehet felhőtároló, Unity Catalog kötet vagy DBFS út. Az első csővezeték frissítéskor a folyamat klónozza az ellenőrzőpontot a csővezeték kezelt tárolójába. Az áramlás ezután az utolsó elkötelezett eltolásról folytatódik, állapota (például aggregációk, deduplikációs kulcsok és vízjelek) érintetlenül. A későbbi csővezeték-frissítések a folyamat klónozott ellenőrzőpontját használják; az eredeti ellenőrzőpontot nem módosítják.
A folyamatnak egy create_table-vel vagy egy elnyelővel létrehozott kezelt táblát kell céloznia.
Állítsd meg az eredeti Strukturált Streaming lekérdezést, mielőtt lefuttatnád a csővezetéket. Az eredeti Strukturált Streaming lekérdezés újrahasználható az import után, de kezelni kell annak ellenőrzési állapotát, és biztosítani kell, hogy a csővezeték és a lekérdezés ne írjon egyszerre ugyanarra a táblára, ami duplikált adatokat eredményezhet.
Hozd létre újra a Strukturált Streaming lekérdezést csővezeték-folyamatként, amely egy új táblába ír, és importálja annak ellenőrzőpontját:
from pyspark import pipelines as dp
# Create a new managed table for the pipeline
dp.create_table("target_table")
# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
target = "target_table",
import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
# The same source your original query read from.
return spark.readStream.table("source_table")
Az ellenőrzőpontot csak egyszer importálják, az első csővezeték-frissítésnél; a későbbi frissítések figyelmen kívül import_checkpointhagyják .
A teljes frissítés nem importálja újra az ellenőrzőpontot; egy új, üres ellenőrzőpontból indul, és újradolgozza a forrást. Egy másik ellenőrzőpont importálásához olyan folyamatnevet használjunk, amelyet korábban nem használtak a céltáblában; egy meglévő folyamatnév újrafelhasználása kihagyja az importot.
Limitations
- Egy ellenőrzőpont importálása egy már létező táblába (például az eredeti Strukturált Streaming lekérdezési célpontba) nem támogatott. Célozz egy új táblát, amit a csővezeték létrehoz, vagy egy elnyelőt.
-
import_checkpointcsak append_flow támogatott.