Ordina l'esecuzione del flusso della pipeline con depends_on

Importante

Questa funzionalità è in Anteprima Pubblica.

Per usare depends_on, configura la tua pipeline per utilizzare il canale PREVIEW di Lakeflow pipelines. Vedere channel in Configurazioni della pipeline.

Di default, una pipeline programma i flussi in base alle dipendenze dei dati: se un flusso legge la tabella che un altro flusso scrive, il lettore corre dietro allo scrittore. L'ordinamento dei flussi permette a un flusso di attendere un altro flusso da cui non legge, dichiarando esplicitamente la dipendenza con depends_on.

depends_on="other_flow" significa che questo flusso inizia solo dopo other_flow essere stato completato con successo. È un vantaggio di programmazione: controlla quando inizia un flusso, non come si svolge. Dichiarare depends_on non cambia il trigger, la modalità o il fatto che sia un flusso a una sola volta. Dichiari comunque questi sul flusso stesso, ad esempio con once=True.

Quando utilizzare l'ordinamento dei flussi

Il caso d'uso principale è drenare dati storici prima di passare a una fonte live mantenendo lo stato di streaming, ad esempio migrando una tabella da un backfill batch a una sorgente Apache Kafka live.

La lettura contemporanea di entrambe le fonti non funziona bene: il backfill limitato trattiene il watermark, ritardando la rimozione dello stato e i risultati basati su finestre. Svuotare prima il buffer e poi avviare la diretta evita questo problema. L'ordinamento del flusso mette i due in sequenza.

L'ordinamento dei flussi supporta anche altri pattern:

  • Ordinare rigorosamente più operazioni di backfill nella stessa tabella.
  • Pipeline in più fasi, come un caricamento iniziale, poi una fase di allineamento, quindi un flusso in tempo reale.
  • I flussi di sequenziamento scrivono su tabelle diverse ma devono essere eseguiti in un ordine prestabilito.

L'ordine fluisce con depends_on

depends_on è disponibile sui decoratori @dp.append_flow e @dp.update_flow nell'API delle pipeline Python. Accetta un unico nome di flusso o una lista di nomi di flusso. Con una lista, ogni flusso nominato deve completarsi prima che inizi il flusso dipendente.

L'esempio seguente convoglia un'operazione di backfill una tantum nella tabella di streaming events e poi avvia un flusso Kafka in tempo reale nella stessa tabella solo dopo il completamento del backfill:

from pyspark import pipelines as dp

dp.create_streaming_table(name="events")

# Drain the historical backfill first.
@dp.append_flow(target="events", once=True, name="events_backfill")
def events_backfill():
    return spark.read.table("historical_events")

# Start the live stream only after the backfill completes.
@dp.append_flow(target="events", name="events_live", depends_on="events_backfill")
def events_live():
    return (
        spark.readStream.format("kafka")
        .option("kafka.bootstrap.servers", "<server>:<port>")
        .option("subscribe", "events")
        .load()
    )

Per attendere più di un predecessore, fornisci un elenco. L'esempio seguente esegue due backfill in parallelo e avvia lo stream live solo dopo che entrambi si completano:

@dp.append_flow(target="events", once=True, name="backfill_2024")
def backfill_2024():
    return spark.read.table("events_2024")

@dp.append_flow(target="events", once=True, name="backfill_2025")
def backfill_2025():
    return spark.read.table("events_2025")

@dp.append_flow(
    target="events",
    name="events_live",
    depends_on=["backfill_2024", "backfill_2025"],
)
def events_live():
    return spark.readStream.format("kafka").option("subscribe", "events").load()

Per eseguire i backfill uno dopo l'altro anziché in parallelo, concatenali con depends_on:

@dp.append_flow(target="events", once=True, name="backfill_2024")
def backfill_2024():
    return spark.read.table("events_2024")

@dp.append_flow(
    target="events", once=True, name="backfill_2025", depends_on="backfill_2024"
)
def backfill_2025():
    return spark.read.table("events_2025")

Requisiti e comportamento

Le seguenti regole si applicano all'ordinamento dei flusso:

  • L'ordinamento dei flussi funziona solo all'interno di una pipeline. Senza il scheduler della pipeline, l'ordine non può essere rispettato.
  • Un predecessore deve scrivere su una tabella o su un sink, non su una vista. Un flusso che scrive in una vista non raggiunge mai uno stato terminale, quindi un flusso ordinato dopo di esso non inizierebbe mai. I target e foreachBatch i sink di tabella di streaming sono entrambi predecessori validi.
  • È consentita l'ordinazione tra destinazioni. Un flusso può dipendere da un flusso che scrive su una tabella diversa.
  • Nomi e cicli di flusso sconosciuti vengono rilevati alla validazione, prima che la pipeline venga eseguita.

Cosa può fungere da predecessore nelle pipeline attivate da eventi e continue

I tipi di flusso che possono fungere da predecessore dipendono dalla modalità di esecuzione della pipeline:

  • Pipeline attivate: qualsiasi flusso può essere un predecessore. Ogni flusso in una corsa attivata raggiunge uno stato terminale, quindi l'ordine si applica all'interno di ogni corsa.
  • Pipeline continue: un predecessore deve essere un flusso monouso (once) che raggiunge uno stato terminale. Un flusso che funziona continuamente non termina mai, quindi un flusso ordinato dopo di esso non inizierebbe mai, e la pipeline lo rifiuta alla validazione.

Poiché un flusso foreachBatch è sempre un sink di streaming e non può essere un flusso monouso, può fungere da predecessore solo in una pipeline attivabile tramite trigger. In una pipeline continua, un foreachBatch flusso può aspettare un once predecessore, ma non può esserlo esso stesso.

Comportamento operativo

Lo stato di completamento di un predecessore once è persistente, quindi si conserva anche dopo i riavvii e gli aggiornamenti della pipeline:

  • Riavvio: i flussi già svuotati restano drenati e vengono ignorati. La pipeline riprende dal primo flusso che non è stato ancora completato. Dopo che l'esecuzione raggiunge il flusso live, i riavvii successivi riprendono solo il flusso live.
  • Aggiornamento completo: cancella lo stato di completamento della catena e lo rifa dall'inizio, in ordine.
  • Ripristino del checkpoint per un singolo flusso: ripristinando il flusso live, viene rieseguito solo il flusso live. I riempimenti a monte once rimangono svuotati e non vengono rieseguiti. Questo è il modo usuale per recuperare una query live senza ri-scaricare la storia.

Risorse aggiuntive