Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Il @dp.append_flow Decorator crea flussi di aggiunta o backfill per le tabelle della pipeline. La funzione deve restituire un dataframe di streaming Apache Spark. Vedere Caricare ed elaborare i dati in modo incrementale con i flussi della pipeline Lakeflow.
Gli add flow possono mirare a tabelle di streaming, tabelle gestite o sink.
Sintassi
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>) #
Parametri
| Parametro | TIPO | Description |
|---|---|---|
| funzione | function |
Obbligatorio. Funzione che restituisce un dataframe di streaming Apache Spark da una query definita dall'utente. |
target |
str |
Obbligatorio. Nome della tabella o del sink che rappresenta la destinazione del flusso di accodamento. |
name |
str |
Nome del flusso. Se non specificato, per impostazione predefinita viene impostato il nome della funzione. |
once |
bool |
Si può facoltativamente definire il flusso come un flusso monouso, ad esempio un riempimento retroattivo. L'uso di once=True modifica il flusso in due modi:
|
depends_on |
str oppure list |
Anteprima pubblica. Uno o più nomi di flusso devono completarsi con successo prima che questo flusso inizi. Accetta un singolo nome di flusso o una lista di nomi. Questo ordina solo l'esecuzione del flusso; non cambia il modo in cui il flusso funziona. Vedi esecuzione del flusso della pipeline di ordini con depends_on. |
comment |
str |
Descrizione del flusso. |
spark_conf |
dict |
Elenco delle configurazioni di Spark per l'esecuzione di questa query |
import_checkpoint |
str |
Il percorso verso un checkpoint di streaming strutturato esistente da importare nel flow, così un flusso migrato riprende dall'ultimo offset commesso invece di rielaborare la sorgente. L'importazione di un checkpoint è in Beta. Vedi Migra un checkpoint di streaming strutturato. |
Esempi
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"))
Migra un checkpoint di streaming strutturato
Importante
L'importazione di un checkpoint è in Beta.
Da usare import_checkpoint per migrare un carico di lavoro di Structured Streaming esistente in una pipeline senza rielaborare il sorgente. Impostalo sulla checkpointLocation tua query di Structured Streaming utilizzata, che può essere uno storage cloud, un volume Unity Catalog o un percorso DBFS. Al primo aggiornamento della pipeline, il flow clona quel checkpoint nello storage gestito della pipeline. Il flusso riprende quindi dall'ultimo offset commesso mantenendo il suo stato (come aggregazioni, chiavi di deduplicazione e filigrane) intatto. Gli aggiornamenti successivi della pipeline utilizzano il checkpoint clonato del flow; il checkpoint originale non viene modificato.
Il flusso deve puntare a una tabella gestita creata con create_table o un pozzo.
Interrompi la query originale di Structured Streaming prima di eseguire la pipeline. La query originale di Structured Streaming può essere riutilizzata dopo l'importazione, ma devi gestire lo stato del checkpoint e assicurarti che la pipeline e la query non scrivano contemporaneamente sulla stessa tabella, il che può produrre dati duplicati.
Ricrea la query di Structured Streaming come un flusso pipeline che scrive in una nuova tabella e importa il suo checkpoint:
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")
Il checkpoint viene importato solo una volta, al primo aggiornamento della pipeline; gli aggiornamenti successivi ignorano import_checkpoint. Un aggiornamento completo non reimporta il checkpoint; parte da un nuovo checkpoint vuoto e riprocessa la sorgente. Per importare un checkpoint diverso, usa un nome di flusso che non sia mai stato usato prima per la tabella di destinazione; riutilizzare un nome di flusso esistente salta l'importazione.
Limitations
- Non è supportato l'importazione di un checkpoint in una tabella già esistente (ad esempio, il target originale di query di Structured Streaming). Punta a una nuova tabella creata dalla pipeline, o a un sink.
-
import_checkpointè supportata solo su append_flow.