append_flow

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:
  • Vrácená hodnota. streaming-query. v tomto případě musí být datový rámec dávky, nikoli datový rámec streamování.
  • Proud je ve výchozím nastavení spuštěn jednou. Pokud je potrubí aktualizováno o úplnou obnovu, ONCE tok se spustí znovu, aby znovu vytvořil data.
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_checkpoint je podporován pouze na append_flow.