Skip to main content

Contrôle de flux

Les bricks de contrôle de flux enchaînent, branchent et sécurisent les exécutions d'un workflow. Elles reçoivent un signal en entrée (in) et l'émettent sur une ou plusieurs sorties selon leur logique. Pour un guide d'usage, voir Orchestrer des pipelines.

Exécuter un pipeline​

Exécute un pipeline de données (sous-pipeline) depuis le workflow.

CatégorieAction
Entréesin (Trigger)
Sortiessuccess (Success) · fail (Fail) · un port par sortie DataFrame du pipeline cible
ParamètreTypeDéfautDescription
pipeline_idchaîne (requis)—UUID du pipeline de données à exécuter.
variablesliste d'objetsvideVariables passées au sous-pipeline (lisibles via ${var.workflow.NOM}). Les valeurs supportent l'interpolation ${...}.
timeout_secondsentier3600Durée maximale. Au-delà, le sous-pipeline et sa descendance sont arrêtés. 0 = sans limite.
output_portsliste d'objetsvideRenseigné automatiquement depuis l'interface du pipeline choisi : un port par sortie DataFrame.

Chaque entrée de variables porte un name (requis) et une value.

Condition​

Branche selon une expression booléenne évaluée sur les variables de workflow ou les résultats des bricks amont. Émet un signal sur la sortie true ou false.

CatégorieAction
Entréesin (Trigger)
Sortiestrue (True) · false (False)
ParamètreTypeDéfautDescription
conditionchaîne (requis)trueExpression Python. Exemple : ${(Nom de l'appel).format} == "xml".

Les références sont remplacées par des littéraux Python : une chaîne se compare donc bien à une chaîne, et un booléen s'utilise seul, sans comparaison. Lisibles ici : ${var.workflow.X}, ${var.item.X} dans une boucle, et ${(Nom de l'appel).sortie} pour le résultat d'un pipeline (voir Interface de pipeline).

Attendre​

Attente conditionnelle dans un workflow : sleep N secondes, attendre une heure donnée, ou attendre un événement (signal webhook).

CatégorieAction
Entréesin (Trigger)
Sortiesout (Resume)
ParamètreTypeDéfautDescription
wait_modeduration · until_time · until_event (requis)durationType d'attente : durée fixe, jusqu'à une heure (HH:MM) ou wait-for-event (signal webhook).

Champs conditionnels selon le mode :

ModeParamètreTypeDéfautDescription
durationduration_secentier (requis)60Durée en secondes (minimum 1).
until_timeuntil_timechaîne (requis)23:00Heure cible au format HH:MM.
until_eventevent_secretchaîne—Secret partagé pour vérifier l'événement entrant (header X-Fluhoms-Event-Secret).
until_eventtimeout_secentier0Échec si aucun événement reçu après ce délai (0 = pas de timeout).

Pour chaque​

Rejoue la sous-pipeline suivante une fois par item. Le corps de boucle est donc une boîte dessinée sur le canvas : sa frontière se voit, au lieu d'être implicite.

CatégorieAction
Entréesin (Trigger)
Sortiesout (Suite)
ParamètreTypeDéfautDescription
source_typepipeline_output · files · variable (requis)pipeline_outputD'où viennent les items.
max_iterationsentier1000Garde-fou. 0 = pas de limite.
stop_on_errorbooléenfalseArrêter au premier item en erreur.
fail_if_emptybooléenfalseÉchouer quand la liste est vide.

Champs conditionnels selon la source :

SourceParamètreDéfautDescription
pipeline_outputoutput_name—Sortie DataFrame d'un pipeline appelé en amont. Vide = la première disponible.
filesconnection_id—Connexion storage à énumérer.
filesfile_pattern*Motif appliqué au nom du fichier.
filesrecursivefalseDescendre dans les sous-dossiers.
variablevariable_value—Tableau JSON, par exemple ${var.workflow.mes_items}.

Lire l'item courant​

RéférenceContenu
${var.item.CHAMP}Un champ de l'item : colonne de la ligne, ou file_name, file_path, file_size… en mode fichiers.
${var.item._index}Rang de l'itération, à partir de 0.
${var.item._total}Nombre total d'items.

L'item vit dans le contexte d'itération de chaque brick, et non dans une variable d'environnement partagée par tout le processus : deux itérations ne peuvent donc pas se marcher dessus.

Comment la boucle s'exécute​

La brick calcule la liste, puis déclare au moteur que la sous-pipeline suivante doit être rejouée. À chaque tour, les bricks de cette sous-pipeline sont réinitialisées et reçoivent le contexte de l'item courant. C'est le même moteur que le mode itératif de la source fichier.

Exécution séquentielle

Les itérations s'enchaînent une par une. Les bricks d'une sous-pipeline sont des instances uniques, réinitialisées entre deux tours : les exécuter en parallèle demanderait de les dupliquer, ce que le moteur ne fait pas encore. Pour paralléliser, le levier disponible est à l'intérieur du traitement — les bricks IA, par exemple, parallélisent déjà leurs appels sur les lignes.

Une seule sous-pipeline en aval

Si plusieurs sous-pipelines suivent la boucle, chacune est parcourue séparément sur la même liste, et un avertissement le signale. Regroupez le corps de boucle dans une seule sous-pipeline.

Parallèle​

Lance simultanément toutes les branches connectées en sortie. Point de départ d'une exécution en éventail (fan-out), à refermer avec une brick Synchronisation.

CatégorieAction
Entréesin (Trigger)
Sortiesout (Branches)
ParamètreTypeDéfautDescription
max_concurrencynombre0Nombre maximum de branches exécutées en parallèle (0 = illimité).

Synchronisation​

Barrière de synchronisation : attend que les branches d'entrée soient terminées avant de poursuivre. Referme un fan-out (Parallèle).

CatégorieAction
Entréesin (Branches)
Sortiesout (Continuer)
ParamètreTypeDéfautDescription
modeall · any · count (requis)allCondition de passage : all = toutes les branches terminées · any = la première terminée · count = un nombre minimum.
countnombre1Nombre de branches à attendre (mode count uniquement).

Réessayer​

Ré-exécute la branche aval en cas d'échec, avec un nombre d'essais et un délai (backoff fixe ou exponentiel).

CatégorieAction
Entréesin (Trigger)
Sortiesout (Succès) · failed (Abandon)
ParamètreTypeDéfautDescription
max_attemptsnombre (requis)3Nombre total de tentatives avant abandon.
backofffixed · exponentialfixedfixed = délai constant · exponential = délai doublé à chaque essai.
delay_secnombre5Délai avant le premier ré-essai, en secondes.

Émet sur out au succès, sur failed après épuisement des essais.

Gestion d'erreur​

Capture l'échec de la branche protégée et bascule sur une branche de repli au lieu d'interrompre le workflow.

CatégorieAction
Entréesin (Flux protégé)
Sortiesout (OK) · on_error (Sur erreur)
ParamètreTypeDéfautDescription
catchall · timeout · assertion (requis)allErreurs capturées : all = toute erreur · timeout = expirations · assertion = échecs d'assertion.

Émet out si tout va bien, on_error en cas d'échec amont.

Assertion​

Vérifie une condition (contrôle qualité). Émet out si elle est vraie. Si elle est fausse : arrête le workflow, route sur failed, ou émet un simple avertissement selon le réglage.

CatégorieAction
Entréesin (Trigger)
Sortiesout (OK) · failed (Échec)
ParamètreTypeDéfautDescription
expressionchaîne (requis)trueExpression Python évaluée. Exemple : ${var.workflow.row_count} > 0.
on_failstop · route · warn (requis)stopstop = stoppe le workflow · route = bascule sur failed · warn = avertit et continue.

Voir aussi​