Ordonner l’exécution des flux du pipeline avec depends_on
Aperçu
Cette fonctionnalité est en aperçu public.
Pour utiliser depends_on, configurez votre pipeline pour qu’il utilise le canal de distribution PREVIEW des Lakeflow pipelines. Consultez channel dans Configurations du pipeline.
By default, un pipeline planifie les flux en fonction de leurs dépendances de données : si un flux lit la table dans laquelle un autre flux écrit, le lecteur s’exécute après l’écrivain. L’ordonnancement des flux permet à un flux d’attendre un autre flux dont il ne lit pas les données, en déclarant explicitement la dépendance avec depends_on.
depends_on="other_flow" signifie que ce flux start uniquement après la fin de l’exécution réussie de other_flow. Il s’agit d’une arête de planification : elle contrôle quand un flux start, et non comment il s’exécute. La déclaration de depends_on ne modifie pas le trigger d’un flux, son mode ou son caractère ponctuel. Vous devez toujours les déclarer sur le flux lui-même, par exemple avec once=True.
Quand utiliser l’ordre des flux
Le cas d’utilisation principal consiste à vider les données historiques avant de basculer vers une source en direct tout en préservant l’état du streaming, comme lors de la migration d’une table issue d’un backfill par batch vers une source Apache Kafka en direct.
La lecture simultanée des deux sources ne fonctionne pas bien : le remplissage borné retarde le filigrane, ce qui retarde l'éviction des états et les résultats fenêtrés. Le vidage préalable du remplissage, suivi du démarrage du live stream, permet d'éviter cela. Le séquençage des flux ordonne les deux.
L’ordre des flux prend également en charge d’autres modèles :
- Ordonnancement strict de plusieurs backfills dans la même table.
- Pipelines par phases, tels qu’un chargement initial, puis une récupération, suivis d’un flux en direct.
- Mise en séquence des flux qui écrivent dans des tables différentes , mais qui doivent s’exécuter dans un ordre précis.
Ordonner les flux avec depends_on
depends_on est disponible sur les décorateurs @dp.append_flow et @dp.update_flow dans l’API Python des pipelines. Il accepte un nom de flux unique ou une liste de noms de flux. Avec une liste, chaque flux nommé doit s’exécuter entièrement avant que le flux dépendant ne start.
L'exemple suivant vide un remplissage unique dans la table de streaming events, puis start un live Kafka Stream dans la même table uniquement après la fin du remplissage :
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()
)
Pour attendre plusieurs prédécesseurs, passez une liste. L’exemple suivant exécute deux remplissages en parallèle et start le live Stream uniquement une fois les deux terminés :
@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()
Pour ordonner les remplissages les uns après les autres plutôt qu’en parallèle, enchaînez depends_on entre eux :
@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")
Conditions requises et comportement
Les règles suivantes s’appliquent à l’ordre des flux :
- L’ordonnancement des flux fonctionne uniquement à l’intérieur d’un pipeline. Sans le planificateur de pipeline, l’ordonnancement ne peut pas être respecté.
- Un prédécesseur doit écrire dans une table ou un récepteur, et non dans une vue. Un flux qui écrit dans une vue n'atteint jamais un état terminal, de sorte qu'un flux ordonné après lui ne start jamais. Les cibles de table de streaming et les récepteurs
foreachBatchsont tous deux des prédécesseurs valides. - Cross-destination ordering is allowed. A flow can depend on a flow that writes to a different table.
- Les noms de flux inconnus et les cycles sont détectés lors de la validation , avant l’exécution du pipeline.
Qu’est-ce qui peut constituer un prédécesseur dans les Trigger et continu pipelines
Les types de flux pouvant agir en tant que prédécesseur dépendent du mode d’exécution du pipeline :
- Pipelines déclenchés : tout flux peut être un prédécesseur. Chaque flux d’une exécution déclenchée atteint un état terminal ; par conséquent, l’ordre s’applique à chaque exécution.
- Pipelines continus : un prédécesseur doit être un flux unique (
once) qui atteint un état terminal. Un flux qui s’exécute en continu ne se termine jamais. Par conséquent, un flux ordonné après celui-ci ne start jamais, et le pipeline le rejette lors de la validation.
Puisqu'un flux foreachBatch est toujours un sink de type streaming et ne peut pas être un flux ponctuel, il ne peut agir en tant que prédécesseur que dans un pipeline déclenché. Dans un pipeline continu, un flux foreachBatch peut attendre un prédécesseur once, mais il ne peut pas lui-même être un prédécesseur.
Comportement opérationnel
L’état d’achèvement d’un prédécesseur once est durable, de sorte qu’il persiste après les redémarrages et les mises à jour du pipeline :
- Redémarrer : les flux déjà drainés restent drainés et sont ignorés. Le pipeline reprend au premier flux qui n’a pas encore été terminé. Une fois que l’exécution a atteint le flux en direct, les redémarrages ultérieurs ne reprennent que le flux en direct.
- Full refresh : clears the completion state for the chain and re-runs it from the beginning, in order.
- Checkpoint reset for a single flow : resetting the live flow replays only the live flow. Les remplissages en amont
oncerestent drainés et ne sont pas réexécutés. Il s’agit de la méthode habituelle pour récupérer une query en direct sans redrainer l’historique.