append_flow

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:
  • Il valore di ritorno. streaming-query. deve essere un dataframe batch in questo caso, non un dataframe di streaming.
  • Il flusso viene eseguito una sola volta per impostazione predefinita. Se la pipeline viene aggiornata con un aggiornamento completo, il ONCE flusso viene eseguito di nuovo per ricreare i dati.
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.