Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
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:
|
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_checkpointwird nur auf append_flow unterstützt.