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.
Dekorátor @dp.append_flow vytvoří přidávací toky nebo backfilly pro tabulky kanálu. Funkce musí vrátit streamingový datový rámec Apache Spark. Viz Načtení a zpracování dat přírůstkově pomocí toků kanálu Lakeflow.
Flow append mohou cílit na streamovací tabulky, spravované tabulky nebo sinky.
Syntaxe
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
depends_on = "<flow-name>", # optional, Public Preview
spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
comment = "<comment>", # optional
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>) #
Parametry
| Parameter | Typ | Description |
|---|---|---|
| funkce | function |
Povinné. Funkce, která vrací datový rámec streamování Apache Sparku z uživatelem definovaného dotazu. |
target |
str |
Povinné. Název tabulky nebo jímky, která je cílem toku připojení. |
name |
str |
Název toku. Pokud není zadaný, nastaví se výchozí hodnota názvu funkce. |
once |
bool |
Volitelně můžete tok definovat jako jednorázový tok, například jako backfill. Použití once=True změní tok dvěma způsoby:
|
depends_on |
str nebo list |
Veřejná ukázka. Jeden nebo více názvů toků, které musí být úspěšně dokončeny, než tento tok začne. Přijímá jedno jméno toku nebo seznam názvů. Toto nařizuje pouze provádění toku; nemění způsob, jakým tok běží. Viz Vykonávání toku toku objednávek s depends_on. |
comment |
str |
Popis toku |
spark_conf |
dict |
Seznam konfigurací Sparku pro spuštění tohoto dotazu |
import_checkpoint |
str |
Cesta k existujícímu kontrolnímu bodu Structured Streaming je potřeba importovat do toku, takže migrovaný stream pokračuje od svého posledního zavázaného offsetu místo opětovného zpracování zdroje. Import checkpointu je v beta fázi. Viz kontrolní bod Migrace strukturovaného streamování. |
Examples
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"))
Migrujte kontrolní bod strukturovaného streamování
Important
Import checkpointu je v beta fázi.
Použijte import_checkpoint k migraci existujícího Structured Streaming workloadu do pipeline bez nutnosti znovu zpracovat zdrojový kód. Nastavte to na checkpointLocation dotaz Structured Streaming, který používáte, což může být cloudové úložiště, svazek Unity Catalog nebo cesta DBFS. Při první aktualizaci pipeline flow tento checkpoint klonuje do spravovaného úložiště pipeline. Tok pak pokračuje od posledního zavázaného posunu se svým stavem (například agregacemi, deduplikačními klíči a vodoznaky) neporušeným. Následné aktualizace pipeline používají klonovaný kontrolní bod toku; původní checkpoint není upravován.
Tok musí cílit na spravovanou tabulku vytvořenou pomocí create_table nebo sinku.
Zastavte původní dotaz na strukturované streamování před spuštěním pipeline. Původní dotaz Structured Streaming lze po importu znovu použít, ale musíte spravovat jeho stav kontrolního bodu a zajistit, aby pipeline a dotaz nezapisovaly do stejné tabulky současně, což může vést k duplicitním datům.
Znovu vytvořte dotaz Structured Streaming jako pipeline flow, který zapisuje do nové tabulky a importuje svůj kontrolní bod:
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")
Kontrolní bod je importován pouze jednou, při první aktualizaci pipeline; pozdější aktualizace ignorují import_checkpoint.
Úplné obnovení kontrolního bodu neimportuje znovu; začíná z nového, prázdného kontrolního bodu a znovu zpracovává zdroj. Pro import jiného kontrolního bodu použijte název toku, který dosud nebyl použit pro cílovou tabulku; opětovné použití existujícího názvu toku přeskočí import.
Omezení
- Import kontrolního bodu do tabulky, která již existuje (například původní cílový dotaz Structured Streaming), není podporován. Zaměřte se na novou tabulku, kterou pipeline vytvoří, nebo na sink.
-
import_checkpointje podporován pouze na append_flow.