append_flow

Der @dp.append_flow Dekorateur erstellt Anfügeflüsse oder Rückfüllungen für Ihre Pipelinetabellen. Die Funktion muss einen Apache Spark Streaming DataFrame zurückgeben. Siehe Load and process data inkremently with Lakeflow pipeline flows.

Append Flows können Streaming-Tabellen, verwaltete Tabellen oder Sinks ansprechen.

Syntax

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

Die Parameter

Parameter Typ Description
Funktion function Erforderlich. Eine Funktion, die einen Apache Spark Streaming DataFrame aus einer benutzerdefinierten Abfrage zurückgibt.
target str Erforderlich. Der Name der Tabelle oder der Senke, die das Ziel des Anfügevorgangs ist.
name str Der Flussname. Wenn nicht angegeben, wird standardmäßig der Funktionsname verwendet.
once bool Definieren Sie optional den Fluss als einmaligen Ablauf, z. B. als Rückfüllvorgang. Durch die Verwendung von once=True wird der Fluss auf zwei Arten verändert:
  • Der Rückgabewert. streaming-query. In diesem Fall muss es sich um einen statischen DataFrame handeln, nicht um einen Streaming-DataFrame.
  • Der Fluss wird standardmäßig einmal ausgeführt. Wenn die Pipeline mit einer vollständigen Aktualisierung aktualisiert wird, wird der ONCE Fluss erneut ausgeführt, um die Daten neu zu erstellen.
depends_on str oder list Öffentliche Vorschau. Ein oder mehrere Flussnamen, die erfolgreich abgeschlossen werden müssen, bevor dieser Fluss beginnt. Akzeptiert einen einzelnen Flussnamen oder eine Liste von Namen. Dies ordnet nur die Ausführung des Flusses an; es ändert nicht, wie der Fluss abläuft. Siehe Order pipeline flow execution with depends_on.
comment str Eine Beschreibung für den Ablauf.
spark_conf dict Eine Liste der Spark-Konfigurationen für die Ausführung dieser Abfrage
import_checkpoint str Der Pfad zu einem bestehenden Structured Streaming-Checkpoint wird in den Flow importiert, sodass ein migrierter Stream von seinem letzten Commitment-Offset fortgesetzt wird, anstatt die Quelle neu zu verarbeiten. Das Importieren eines Kontrollpunkts befindet sich in der Beta. Siehe Migrate a Structured Streaming Checkpoint.

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

Migrieren eines strukturierten Streaming-Checkpoints

Important

Das Importieren eines Kontrollpunkts befindet sich in der Beta.

Nutze import_checkpoint es, um eine bestehende Structured Streaming-Workload in eine Pipeline zu migrieren, ohne den Quellcode neu zu verarbeiten. Setze es auf die checkpointLocation verwendete Structured Streaming-Abfrage, die ein Cloud-Speicher, ein Unity-Catalog-Volume oder ein DBFS-Pfad sein kann. Beim ersten Pipeline-Update klont der Flow diesen Checkpoint in den verwalteten Speicher der Pipeline. Der Fluss wird dann mit seinem Zustand (wie Aggregationen, Deduplizierungsschlüsseln und Wasserzeichen) vom zuletzt zugesagten Offset fortgesetzt. Nachfolgende Pipeline-Aktualisierungen verwenden den geklonten Checkpoint des Flows; der ursprüngliche Checkpoint wird nicht verändert.

Der Flow muss eine verwaltete Tabelle ansprechen, die mit create_table oder einem Sink erstellt wurde.

Stoppen Sie die ursprüngliche Structured Streaming-Abfrage, bevor Sie die Pipeline ausführen. Die ursprüngliche Structured Streaming-Abfrage kann nach dem Import wiederverwendet werden, aber Sie müssen ihren Checkpoint-Status verwalten und sicherstellen, dass Pipeline und Abfrage nicht gleichzeitig in dieselbe Tabelle schreiben, was doppelte Daten erzeugen kann.

Erstellen Sie die Structured Streaming-Abfrage als Pipeline-Fluss, der in eine neue Tabelle schreibt und deren Checkpoint importiert:

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

Der Checkpoint wird nur einmal importiert, beim ersten Pipeline-Update; spätere Updates ignorieren import_checkpoint. Eine vollständige Aktualisierung importiert den Checkpoint nicht erneut; er startet von einem neuen, leeren Checkpoint und verarbeitet den Quellcode erneut. Um einen anderen Checkpoint zu importieren, verwenden Sie einen Flussnamen, der zuvor für die Zieltabelle nicht verwendet wurde; die Wiederverwendung eines bestehenden Flussnamens überspringt den Import.

Einschränkungen

  • Das Importieren eines Checkpoints in eine bereits existierende Tabelle (zum Beispiel das ursprüngliche Structured Streaming Abfrageziel) wird nicht unterstützt. Ziele eine neue Tabelle an, die die Pipeline erstellt, oder einen Sink.
  • import_checkpoint wird nur auf append_flow unterstützt.