Modes d’exécution du flux de travail

Lors de l’exécution d’un flux de travail dans .NET, le mode d’exécution contrôle la façon dont les supersteps sont traités et la façon dont les événements sont remis au consommateur. La InProcessExecution classe expose deux modes d’exécution : OffThread et Lockstep.

Overview

OffThread (valeur par défaut) Lockstep
Exécution de superstep Thread d’arrière-plan Thread du consommateur
Livraison d’événements Immédiatement, au fur et à mesure que les événements sont déclenchés Regroupés après la fin de chaque super-étape
Exécution de l’étape Indépendant du traitement des événements Mis en pause jusqu’à ce que les événements regroupés soient consommés
Concurrence Le consommateur lit les événements pendant l’exécution des super-étapes Le consommateur et l’exécution des super-étapes alternent
Idéal pour Streaming en temps réel, scénarios de production Test, débogage, classement déterministe

En dehors du thread

OffThread est le mode d’exécution par défaut . Les super-étapes s’exécutent sur un thread d’arrière-plan, et les événements sont diffusés immédiatement au moment où ils sont générés via une implémentation basée sur des canaux.

// OffThread is the default — these are equivalent:
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input);
await using StreamingRun run = await InProcessExecution.OffThread.RunStreamingAsync(workflow, input);

Fonctionnement

  1. Une tâche en arrière-plan exécute en continu les super-étapes tant que des messages sont en attente.
  2. À mesure que les exécuteurs produisent des sorties ou des événements, les objets résultants WorkflowEvent sont écrits dans un objet non lié Channel<WorkflowEvent>.
  3. Le consommateur lit les événements à partir du canal via WatchStreamAsync, les recevant en temps réel à mesure qu’ils sont produits.
  4. Lorsque tous les supersteps sont terminés et qu'aucun message ne reste à traiter, l'exécution s'arrête avec un statut Idle ou PendingRequests.

Comme la boucle des super-étapes et le consommateur s’exécutent en parallèle, les événements apparaissent dès qu’ils sont générés : il n’existe aucun délai lié à une mise en mémoire tampon. Cela rend le mode OffThread particulièrement adapté aux scénarios de diffusion en continu (streaming) nécessitant une faible latence de transmission des événements, par exemple l’affichage de mises à jour jeton par jeton dans une interface utilisateur.

Exécutions simultanées

OffThread prend également en charge une variante simultanée qui permet à plusieurs exécutions de partager simultanément la même instance de flux de travail :

await using StreamingRun run = await InProcessExecution.Concurrent.RunStreamingAsync(workflow, input);

Important

L’exécution simultanée nécessite que tous les exécuteurs du flux de travail soient déclarés crossRunShareable (sur le constructeur) ou fournis en tant que méthodes de fabrique.

Lockstep

En mode Lockstep, les super-étapes s’exécutent dans le thread du consommateur plutôt que dans une tâche en arrière-plan. Les événements sont accumulés pendant chaque superstep et émis en tant que lot une fois le superstep terminé.

await using StreamingRun run = await InProcessExecution.Lockstep.RunStreamingAsync(workflow, input);

Fonctionnement

  1. Le consommateur appelle WatchStreamAsync, qui pilote la boucle d’exécution.
  2. Une super-étape s’exécute jusqu’à son terme, et les événements sont accumulés dans une file d’attente.
  3. Une fois la super-étape terminée, tous les événements en attente dans la file sont transmis au consommateur.
  4. La super-étape suivante ne démarre qu’après que le consommateur a reçu tous les événements de la précédente.

Ce modèle alternatif signifie que le consommateur et le moteur de workflow ne s’exécutent jamais simultanément. La livraison des événements est déterministe : tous les événements d’une super-étape sont garantis d’arriver avant tout événement de la super-étape suivante.

Quand utiliser Lockstep

Lockstep est utile lorsque :

  • Test : l’ordre des événements déterministe rend les assertions simples.
  • Débogage : le débogage pas à pas est plus simple lorsque l’exécution reste dans le thread du consommateur.
  • Traitement ordonné : scénarios dans lesquels vous devez traiter entièrement les événements d’un superstep avant le début du prochain superstep.

Choix d’un mode d’exécution

Pour la plupart des scénarios de production, le mode OffThread par défaut est recommandé. Il offre la meilleure réactivité et permet au flux de travail de poursuivre le traitement pendant que le consommateur gère les événements.

Utilisez Lockstep lorsque le comportement déterministe est plus important que les performances, comme dans les tests unitaires ou les sessions de débogage.

// Production: OffThread (default)
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input);

// Testing: Lockstep for deterministic behavior
await using StreamingRun run = await InProcessExecution.Lockstep.RunStreamingAsync(workflow, input);

Exécution sans diffusion

Les deux modes d’exécution prennent en charge l’exécution non continue via RunAsync. En mode non streaming, le flux de travail s’exécute jusqu’à la fin et collecte tous les événements dans un Run objet plutôt que de les diffuser en continu de manière incrémentielle :

Run run = await InProcessExecution.RunAsync(workflow, input);

// Access all emitted events
foreach (WorkflowEvent evt in run.OutgoingEvents)
{
    // Process events
}

Étant donné que l’exécution hors streaming collecte tous les événements après la fin, l’avantage de remise d’événements en temps réel de OffThread ne s’applique pas. La principale différence entre les modes dans les scénarios de non-diffusion en continu est threading : OffThread exécute des supersteps sur un thread d’arrière-plan, libérant le thread appelant en attendant l’achèvement, tandis que Lockstep exécute des supersteps sur le thread de l’appelant, le bloquant jusqu’à ce que le flux de travail se termine.

L’exécution hors streaming utilise le mode OffThread par défaut. Pour utiliser Lockstep avec une exécution sans diffusion en continu :

Run run = await InProcessExecution.Lockstep.RunAsync(workflow, input);

Étapes suivantes

Les modes d’exécution ne s’appliquent pas aux flux de travail Python. Les flux de travail Python utilisent un modèle d’exécution unique qui gère le traitement et la distribution d’événements en superstep via un générateur asynchrone. Ce modèle est similaire au mode Lockstep .NET. Les étapes ne sont pas avancées, sauf si le consommateur extrait activement les événements du générateur.

Pour plus d’informations sur l’exécution de flux de travail Python, consultez Workflow Builder &Execution.

Lors de l’exécution d’un flux de travail dans Go, l’environnement d’exécution contrôle la façon dont les supersteps sont traités et la façon dont les événements sont remis au consommateur. Le workflow/inproc package expose trois environnements : Default/OffThread, Lockstepet .Concurrent

Overview

OffThread / Par défaut Lockstep Simultané
Exécution de superstep Goroutine d’arrière-plan Pilotée par le consommateur d’événements Goroutine d’arrière-plan
Livraison d’événements Immédiatement, au fur et à mesure que les événements sont déclenchés Regroupée au fur et à mesure de la consommation du flux Immédiatement, au fur et à mesure que les événements sont déclenchés
Idéal pour Streaming en temps réel, scénarios de production Test, débogage, classement déterministe Instances de flux de travail partagées avec liaisons sécurisées simultanées

En dehors du thread

OffThread est le mode d’exécution par défaut. Ceux-ci sont équivalents :

stream, err := inproc.Default.RunStreaming(ctx, wf, input)
stream, err := inproc.OffThread.RunStreaming(ctx, wf, input)

Fonctionnement

  1. Une goroutine en arrière-plan exécute les super-étapes tant que des messages sont en attente.
  2. À mesure que les exécuteurs produisent des sorties ou des événements, les événements de flux de travail sont écrits dans le flux.
  3. Le consommateur lit les événements à l’aide de WatchStream, en les recevant au fur et à mesure de leur production.
  4. Lorsque toutes les super-étapes sont terminées et qu’il ne reste aucun message, l’exécution s’arrête avec un état inactif ou demande en attente.

Exécutions simultanées

Utilisez inproc.Concurrent quand toutes les liaisons d’exécuteur dans le flux de travail prennent en charge l’exécution partagée simultanée :

stream, err := inproc.Concurrent.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer stream.Close(ctx)

Lockstep

En mode Lockstep, l’exécution du workflow progresse au fur et à mesure que le consommateur lit le flux. Cela rend l’ordre des événements déterministe pour les tests et le débogage.

stream, err := inproc.Lockstep.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer stream.Close(ctx)

for evt, err := range stream.WatchStream(ctx) {
    if err != nil {
        return err
    }
    // inspect event
}

Fonctionnement

  1. Le consommateur appelle WatchStream, qui pilote la boucle d’exécution.
  2. Une super-étape s’exécute jusqu’à son terme, et les événements sont accumulés.
  3. Les événements accumulés sont transmis au consommateur.
  4. La super-étape suivante ne démarre qu’après que le consommateur a reçu les événements de la super-étape précédente.

Quand utiliser Lockstep

Utilisez Lockstep lorsque le comportement déterministe importe plus que la diffusion en continu à faible latence, comme les tests unitaires, le débogage ou les scénarios dans lesquels vous souhaitez traiter entièrement les événements d’un superstep avant le début du superstep suivant.

Choix d’un mode d’exécution

Pour la plupart des scénarios de production, utilisez inproc.Default ou inproc.OffThread. Utiliser inproc.Lockstep quand l’ordre des événements déterministe est plus important que la latence de diffusion en continu, comme dans les tests. Utilisez inproc.Concurrent uniquement lorsque chaque liaison du flux de travail prend en charge l’exécution partagée simultanée.

// Production: OffThread (default)
stream, err := inproc.Default.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer stream.Close(ctx)

// Testing: Lockstep for deterministic behavior
testStream, err := inproc.Lockstep.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer testStream.Close(ctx)

Exécution sans diffusion

Tous les environnements d’exécution prennent également en charge l’exécution sans diffusion (non-streaming) via Run. Celle-ci s’exécute jusqu’au prochain point d’arrêt, puis stocke les événements générés dans l’exécution retournée.

run, err := inproc.Default.Run(ctx, wf, input)
if err != nil {
    return err
}

for evt := range run.NewEvents() {
    if output, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Final result: %v\n", output.Output)
    }
}

Étapes suivantes