Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Az al-munkafolyamatok olyan teljes munkafolyamatok, amelyek végrehajtóként futnak egy szülő-munkafolyamaton belül. Ez lehetővé teszi összetett rendszerek írását kisebb, újrafelhasználható munkafolyamat-építőelemekből – mindegyik saját elkülönített végrehajtási környezettel, állapotkezeléssel és üzenet-útválasztással.
Overview
Az al-munkafolyamatok akkor hasznosak, ha:
- Összetettség felbontása – egy nagy munkafolyamat kisebb, egymástól függetlenül tesztelhető egységekre bontható.
- Munkafolyamat-logika újrafelhasználása – ugyanazt az al-munkafolyamatot ágyazza be több szülő-munkafolyamatba.
- Az állapot elkülönítése – az egyes al-munkafolyamatok belső állapotának elkülönítése a szülőtől.
- Az adatfolyam szabályozása – az üzenetek csak az élein keresztül lépnek be és hagyják el az al-munkafolyamatot, és nem történik közvetítés több szinten keresztül.
Amikor egy al-munkafolyamatot hozzáad egy szülő munkafolyamathoz, ugyanúgy viselkedik, mint bármely más végrehajtó: bemeneti üzeneteket fogad, belső gráfját a befejezésig futtatja, és kimeneti üzeneteket hoz létre az alsóbb rétegbeli végrehajtók számára.
Sub-Workflow létrehozása
A C#-ban az al-munkafolyamatok kétféleképpen írhatóak:
-
Közvetlen kötés – munkafolyamat
BindAsExecutor()közvetlen beágyazása végrehajtóként a szülő munkafolyamatba. Ez megőrzi az al-munkafolyamat natív bemeneti/kimeneti típusait. -
Ügynökburkolás – használja a
AsAIAgent()-t a munkafolyamat ügynökké alakításához, majd adja hozzá az ügynököt a szülő munkafolyamathoz. Ez akkor hasznos, ha a szülő munkafolyamat ügynökalapú végrehajtókat használ.
Közvetlen kötés a BindAsExecutorral
A BindAsExecutor() bővítménymetódus átalakít egy munkafolyamatot olyan ExecutorBinding munkafolyamattá, amely közvetlenül hozzáadható egy szülő munkafolyamathoz:
using Microsoft.Agents.AI.Workflows;
// Create executors for the inner workflow
UppercaseExecutor uppercase = new();
ReverseExecutor reverse = new();
AppendSuffixExecutor append = new(" [PROCESSED]");
// Build the inner workflow
var innerWorkflow = new WorkflowBuilder(uppercase)
.AddEdge(uppercase, reverse)
.AddEdge(reverse, append)
.WithOutputFrom(append)
.Build();
// Bind the inner workflow as an executor
ExecutorBinding subWorkflowExecutor = innerWorkflow.BindAsExecutor("TextProcessingSubWorkflow");
// Build the parent workflow using the sub-workflow executor
PrefixExecutor prefix = new("INPUT: ");
PostProcessExecutor postProcess = new();
var parentWorkflow = new WorkflowBuilder(prefix)
.AddEdge(prefix, subWorkflowExecutor)
.AddEdge(subWorkflowExecutor, postProcess)
.WithOutputFrom(postProcess)
.Build();
Ezzel BindAsExecutoraz al-munkafolyamat beírt bemeneti és kimeneti típusai megmaradnak – a szülő munkafolyamat az al-munkafolyamat által várt és előállított tényleges típusok alapján irányítja át az üzeneteket.
Ügynökburkolás az AsAIAgenttel
Amikor a szülő munkafolyamat ügynökalapú végrehajtókat használ, konvertálja a belső munkafolyamatot ügynökhasználatra a következő módon: AsAIAgent(). Az WorkflowBuilder ügynök automatikusan egy végrehajtóba burkolódik:
using Microsoft.Agents.AI;
using Microsoft.Agents.AI.Workflows;
// Create agents for the inner workflow
AIAgent specialist1 = chatClient.AsAIAgent("You are specialist 1. Analyze the data.");
AIAgent specialist2 = chatClient.AsAIAgent("You are specialist 2. Validate the analysis.");
// Build the inner workflow
var innerWorkflow = new WorkflowBuilder(specialist1)
.AddEdge(specialist1, specialist2)
.Build();
// Convert the inner workflow to an agent
AIAgent innerWorkflowAgent = innerWorkflow.AsAIAgent(
id: "analysis-pipeline",
name: "Analysis Pipeline",
description: "A sub-workflow that analyzes and validates data"
);
// Create agents for the parent workflow
AIAgent coordinator = chatClient.AsAIAgent("You are a coordinator. Delegate tasks to the team.");
AIAgent reviewer = chatClient.AsAIAgent("You are a reviewer. Review the final output.");
// Build the parent workflow with the sub-workflow
var parentWorkflow = new WorkflowBuilder(coordinator)
.AddEdge(coordinator, innerWorkflowAgent)
.AddEdge(innerWorkflowAgent, reviewer)
.Build();
A belső munkafolyamat egyetlen lépésben fut a szülő munkafolyamat szemszögéből. A koordinátor üzeneteket küld az elemzési folyamatnak, amely belsőleg fut specialist1 → specialist2, majd továbbítja az eredményt a véleményezőnek.
Tip
Használja a BindAsExecutor() elemet gépelt végrehajtókkal és a AsAIAgent() elemet ügynökalapú munkafolyamatokkal való munka során. A munkafolyamatok ügynökké alakításának konfigurálásáról további információt a Munkafolyamatok ügynökként című témakörben talál.
Bemeneti és kimeneti típusok
Ha egy munkafolyamatot al-munkafolyamatként használ, megőrzi a belső végrehajtók típusszerződéseit.
Ezzel BindAsExecutoregyütt az al-munkafolyamat-végrehajtó ugyanazokat a bemeneti típusokat fogadja el, mint a belső munkafolyamat kezdő végrehajtója, és ugyanazokat a kimeneti típusokat küldi el, amelyeket a belső munkafolyamat előállít. A szülő munkafolyamat éleinek össze kell kapcsolniuk azokat a végrehajtókat, amelyek kimeneti típusai megegyeznek az al-munkafolyamat várt bemeneti típusaival, és az al-munkafolyamat kimeneti típusainak meg kell egyeznie az alsóbb rétegbeli végrehajtók várt bemenetével.
A AsAIAgent segítségével az al-munkafolyamat ügynökként van becsomagolva, és követi az Ügynök Végrehajtó bemeneti/kimeneti szerződéseit (string, ChatMessage, IEnumerable<ChatMessage>).
Kimeneti viselkedés
Alapértelmezés szerint, amikor egy al-munkafolyamat kimeneteket állít elő (a YieldOutputAsync segítségével), azokat üzenetekként továbbítják a szülő munkafolyamat csatlakoztatott végrehajtóinak. Ez lehetővé teszi az alsóbb rétegbeli végrehajtók számára a munkafolyamat-részeredmények feldolgozását.
Az ExecutorOptions osztály a következő viselkedést vezérli:
| Option | Alapértelmezett | Description |
|---|---|---|
AutoSendMessageHandlerResultObject |
true |
Az al-munkafolyamat kimeneteinek továbbítása üzenetként a csatlakoztatott végrehajtóknak a szülőgráfban. |
AutoYieldOutputHandlerResultObject |
false |
Az al-munkafolyamat kimeneteit közvetlenül a szülő munkafolyamat kimeneti eseménystreamjének adja vissza. |
Ha AutoYieldOutputHandlerResultObject engedélyezve van, az al-munkafolyamat kimenetei megkerülik a szülő belső útválasztását, és közvetlenül a szülő munkafolyamat hívójának lesznek kézbesítve.
var options = new ExecutorOptions
{
AutoYieldOutputHandlerResultObject = true,
};
ExecutorBinding subWorkflowExecutor = innerWorkflow.BindAsExecutor("SubWorkflow", options);
Kérelmek és válaszok
Az al-munkafolyamatok teljes mértékben támogatják a kérés- és válaszmechanizmust . Amikor az al-munkafolyamatban egy végrehajtó kérést küld (például emberi bemenet igénylésére), a WorkflowHostExecutor rendszer egy RequestInfoEvent továbbítja a szülő munkafolyamatnak – az al-munkafolyamat végrehajtójának azonosítója elő van állítva a portazonosítóra (példáulSubWorkflow.GuessNumber).
Ez a minősítés biztosítja, hogy amikor a szülő munkafolyamat választ kap, vissza tudja irányítani a választ a megfelelő al-munkafolyamat-példányra. A szülő munkafolyamat az al-munkafolyamat-kérelmeket ugyanazzal a válaszmechanizmussal kezeli, mint bármely más kérés:
await using StreamingRun handle = await InProcessExecution.RunStreamingAsync(parentWorkflow, input);
await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
switch (evt)
{
case RequestInfoEvent requestInfoEvt:
// The request may originate from the sub-workflow
// Handle it and send the response back
var response = requestInfoEvt.Request.CreateResponse(myResponseData);
await handle.SendResponseAsync(response);
break;
case WorkflowOutputEvent outputEvt:
Console.WriteLine($"Output: {outputEvt.Data}");
break;
}
}
Megjegyzés:
A fölérendelt munkafolyamat-hívó szempontjából nincs különbség a legfelső szintű végrehajtótól érkező kérések és az al-munkafolyamatból érkező kérések között. A keretrendszer transzparensen kezeli az útválasztást.
Hogyan működik?
Amikor a szülő munkafolyamat átirányít egy üzenetet az al-munkafolyamat végrehajtójának:
-
Bemeneti kézbesítés – az üzenet a belső munkafolyamat kezdő végrehajtójának továbbítja. Ezzel
BindAsExecutoraz üzenettípusnak meg kell egyeznie a kezdő végrehajtó várt típusával. AzAsAIAgenthasználatával az üzenetekChatMessageformátumra normalizálódnak. - Belső végrehajtás – a belső munkafolyamat saját superstep ciklust futtat.
-
Kimeneti gyűjtemény – a rendszer összegyűjti a belső munkafolyamat kimeneti eseményeit. A
BindAsExecutorkimenetek megtartják az eredeti típusukat. AAsAIAgentkimenetek ügynöki válaszüzenetekké alakulnak. - Kérelemtovábbítás – ha a belső munkafolyamat függőben lévő kérésekkel rendelkezik, a rendszer továbbítja őket a szülő munkafolyamatnak kezelés céljából (lásd : Kérések és válaszok).
- Alsóbb rétegbeli küldés – az így kapott üzenetek a szülő munkafolyamat következő végrehajtójának lesznek elküldve.
Mivel a belső munkafolyamat saját végrehajtási környezetet tart fenn, az állapota független a szülő-munkafolyamattól.
Tip
A munkafolyamatok ügynökké alakításának konfigurálásáról, beleértve a streamelési viselkedést és a kivételkezelést, tekintse meg a Munkafolyamatok ügynökként című témakört.
Többszintű beágyazás
Az al-munkafolyamatok tetszőleges mélységbe ágyazhatók. Minden szint saját végrehajtási környezetet tart fenn:
// Level 1: Data preparation pipeline
var dataPipeline = new WorkflowBuilder(fetcher)
.AddEdge(fetcher, cleaner)
.Build();
AIAgent dataPipelineAgent = dataPipeline.AsAIAgent(
id: "data-pipeline",
name: "Data Pipeline"
);
// Level 2: Analysis pipeline (contains the data pipeline)
var analysisPipeline = new WorkflowBuilder(dataPipelineAgent)
.AddEdge(dataPipelineAgent, analyzer)
.Build();
AIAgent analysisPipelineAgent = analysisPipeline.AsAIAgent(
id: "analysis-pipeline",
name: "Analysis Pipeline"
);
// Level 3: Top-level orchestration
var topWorkflow = new WorkflowBuilder(coordinator)
.AddEdge(coordinator, analysisPipelineAgent)
.AddEdge(analysisPipelineAgent, reporter)
.Build();
Megjegyzés:
Minden beágyazási szint végrehajtási többletterhelést eredményez, mivel a belső munkafolyamat saját szupersztep-ciklust futtat. Tartsa ésszerű határok között a fészkelés mélységét a teljesítményérzékeny forgatókönyvekben.
Hibakezelés
Ha egy al-munkafolyamat meghiúsul, a hiba propagálása a szülő munkafolyamatba SubworkflowErrorEventtörténik. A szülő munkafolyamat az eseménystreamen keresztül figyelheti meg ezeket a hibákat:
await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
if (evt is SubworkflowErrorEvent subError)
{
Console.WriteLine($"Sub-workflow '{subError.ExecutorId}' failed: {subError.Data}");
}
}
Ha az al-munkafolyamat kezeletlen kivételt tapasztal, a szülő munkafolyamat végrehajtása folytatódik, de az al-munkafolyamat végrehajtója leállítja a további üzenetek feldolgozását.
Ellenőrző pontok használata
Amikor ellenőrzőpontot hoz létre a szülő-munkafolyamaton, a rendszer szerializálja az al-munkafolyamat-ügynök munkamenet-állapotát a szülő végrehajtó ellenőrzőpont-adatainak részeként. A visszaállítás során a munkamenet állapota deszerializálódik, lehetővé téve a szülő munkafolyamat számára, hogy a változatlan alfolyamat állapotával folytatódjon.
CheckpointManager checkpointManager = CheckpointManager.CreateInMemory();
// Run the parent workflow with checkpointing
StreamingRun run = await InProcessExecution
.RunStreamingAsync(parentWorkflow, input, checkpointManager);
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
// Process events, including those from sub-workflows
}
// Resume from a checkpoint
CheckpointInfo checkpoint = run.Checkpoints[^1];
StreamingRun resumedRun = await InProcessExecution
.ResumeStreamingAsync(parentWorkflow, checkpoint, checkpointManager);
Sub-Workflow létrehozása
Pythonban egy al-munkafolyamatot úgy hozhat létre, hogy egy Workflow beágyaz egy WorkflowExecutor-be, és hozzáadja azt egy szülő munkafolyamathoz.
from agent_framework import WorkflowBuilder, WorkflowExecutor
# Create agents for the inner workflow
specialist1 = client.as_agent(name="Specialist1", instructions="Analyze the data.")
specialist2 = client.as_agent(name="Specialist2", instructions="Validate the analysis.")
# Build the inner workflow
inner_workflow = (
WorkflowBuilder(start_executor=specialist1)
.add_edge(specialist1, specialist2)
.build()
)
# Wrap as an executor
inner_workflow_executor = WorkflowExecutor(
workflow=inner_workflow,
id="analysis-pipeline",
)
# Create agents for the parent workflow
coordinator = client.as_agent(name="Coordinator", instructions="Delegate tasks to the team.")
reviewer = client.as_agent(name="Reviewer", instructions="Review the final output.")
# Build the parent workflow with the sub-workflow
parent_workflow = (
WorkflowBuilder(start_executor=coordinator)
.add_edge(coordinator, inner_workflow_executor)
.add_edge(inner_workflow_executor, reviewer)
.build()
)
A belső munkafolyamat egyetlen lépésben fut a szülő munkafolyamat szemszögéből. A koordinátor üzeneteket küld az elemzési folyamatnak, amely belsőleg fut specialist1 → specialist2, majd továbbítja az eredményt a véleményezőnek.
WorkflowExecutor paraméterek
| Paraméter | Típus | Alapértelmezett | Description |
|---|---|---|---|
workflow |
Workflow |
— | A munkafolyamat-példány, amely végrehajtóként van körbefuttatva. |
id |
str |
— | A végrehajtó egyedi azonosítója. |
allow_direct_output |
bool |
False |
Amikor Trueaz al-munkafolyamat kimenetei közvetlenül a szülő munkafolyamat eseményfolyamára kerülnek, ahelyett, hogy üzenetként küldenék őket a csatlakoztatott végrehajtóknak. |
propagate_request |
bool |
False |
Amikor az al-munkafolyamatból érkező kérések a szülő munkafolyamat eseményfolyamára kerülnek propagálásra, rendszeres kérésinformációs eseményekként jelennek meg. Amikor Falsea rendszer becsomagolta SubWorkflowRequestMessage a kérelmeket a fölérendelt végrehajtók elfogására. |
Almunkafolyamatok beágyazása
A Workflow példányait kifejezetten foglalja egy WorkflowExecutor elembe, mielőtt hozzáadja őket egy szülő munkafolyamathoz. Az ügynökök közvetlenül átadhatók a(z) WorkflowBuilder számára, de a nyers Workflow példányokhoz ez a burkoló szükséges.
from agent_framework import WorkflowExecutor
inner_workflow_executor = WorkflowExecutor(inner_workflow, id="analysis_pipeline")
parent_workflow = (
WorkflowBuilder(start_executor=coordinator)
.add_edge(coordinator, inner_workflow_executor)
.add_edge(inner_workflow_executor, reviewer)
.build()
)
Az explicit körbefuttatás lehetővé teszi a következőket:
- Rendeljen hozzá egy adott végrehajtóazonosítót a hivatkozáshoz több élen.
- Használja ugyanazt a
WorkflowExecutorpéldányt újra a gráfon.
# Explicit wrapping — create the WorkflowExecutor yourself
inner_workflow_executor = WorkflowExecutor(
workflow=inner_workflow,
id="analysis-pipeline",
)
parent_workflow = (
WorkflowBuilder(start_executor=coordinator)
.add_edge(coordinator, inner_workflow_executor)
.add_edge(inner_workflow_executor, reviewer)
.build()
)
Bemeneti és kimeneti típusok
A becsomagolt munkafolyamat örökli a típusjelölését:
-
A bemeneti típusok megfelelnek a becsomagolt munkafolyamat kezdő végrehajtójának bemeneti típusainak (valamint
SubWorkflowResponseMessagea továbbított kérelmekre adott válaszok kezeléséhez). -
A kimeneti típusok megegyeznek a becsomagolt munkafolyamat kimeneti típusaival. Ha az al-munkafolyamat bármely végrehajtója képes a kérés-válaszra,
SubWorkflowRequestMessageakkor kimeneti típusként is szerepel.
Ez azt jelenti, hogy a szülő munkafolyamat éleinek össze kell kapcsolniuk azokat a végrehajtókat, amelyek kimeneti típusai megegyeznek az al-munkafolyamat várt bemeneti típusaival. Hasonlóképpen, az alsóbb rétegbeli végrehajtóknak el kell fogadniuk az al-munkafolyamat által előállított típusokat:
# The sub-workflow's start executor accepts TextProcessingRequest
# So the parent executor must send TextProcessingRequest
class Orchestrator(Executor):
@handler
async def start(self, texts: list[str], ctx: WorkflowContext[TextProcessingRequest]) -> None:
for text in texts:
await ctx.send_message(TextProcessingRequest(text=text))
# The sub-workflow yields TextProcessingResult
# So the downstream executor must handle TextProcessingResult
class ResultCollector(Executor):
@handler
async def collect(self, result: TextProcessingResult, ctx: WorkflowContext) -> None:
print(f"Received: {result}")
Kimeneti viselkedés
Alapértelmezés szerint (allow_direct_output=False) amikor egy al-munkafolyamat kimeneteket hoz létre, yield_outputa rendszer üzenetekként továbbítja ezeket a kimeneteket a szülő munkafolyamatban send_messagelévő csatlakoztatott végrehajtóknak. Ez lehetővé teszi az alsóbb rétegbeli végrehajtók számára, hogy a fölérendelt gráf részeként dolgozzák fel az al-munkafolyamat eredményeit.
Amikor allow_direct_output=True, az al-munkafolyamat kimenetei közvetlenül a szülő munkafolyamat eseményfolyamához kerülnek. Az al-munkafolyamat kimenetei a szülő munkafolyamat kimenetei lesznek, megkerülve a szülő belső végrehajtói útválasztását:
# Outputs go directly to parent's event stream
sub_workflow_executor = WorkflowExecutor(
workflow=inner_workflow,
id="analysis-pipeline",
allow_direct_output=True,
)
# The caller receives sub-workflow outputs directly
async for event in parent_workflow.run(input_data, stream=True):
if event.type == "output":
# This output came from the sub-workflow
print(event.data)
Köztes kibocsátások gyermek-munkafolyamatokból
"intermediate" a gyermek-munkafolyamatban létrehozott események automatikusan továbbítódnak a szülő munkafolyamat eseményfolyamába. Ezek a WorkflowExecutor saját id elemének tulajdoníthatók (nem annak a belső végrehajtónak, amely eredetileg kibocsátotta őket), ami megőrzi az enkapszulációt. Döntő fontosságú, hogy ezek az események megtartják a "intermediate" címkét , függetlenül attól, hogy a szülő hogyan jelöli meg a WorkflowExecutor saját output_from vagy intermediate_output_from listáiban.
async for event in parent_workflow.run(input_data, stream=True):
if event.type == "intermediate":
# Attributed to the WorkflowExecutor id, e.g. "analysis-pipeline"
print(f"[{event.executor_id}] intermediate: {event.data}")
elif event.type == "output":
print(f"Terminal output: {event.data}")
Kérelmek és válaszok
Az al-munkafolyamatok teljes mértékben támogatják a kérés- és válaszmechanizmust . Amikor egy végrehajtó egy al-munkafolyamatban hív ctx.request_info(), a WorkflowExecutor rendszer elfogja a kérést, és a propagate_request beállítás alapján kezeli azt.
Kérelmek lehallgatása a szülői munkafolyamatban (alapértelmezett)
Alapértelmezés szerint propagate_request=False, az al-munkafolyamatból érkező kérések be vannak csomagolva SubWorkflowRequestMessage-be, és a szülő munkafolyamat csatlakoztatott végrehajtóinak lesznek elküldve. Ez lehetővé teszi, hogy a szülő végrehajtók helyileg kezeljék a kérést:
from agent_framework import (
SubWorkflowRequestMessage,
SubWorkflowResponseMessage,
)
class ParentHandler(Executor):
@handler
async def handle_request(
self,
request: SubWorkflowRequestMessage,
ctx: WorkflowContext[SubWorkflowResponseMessage],
) -> None:
# Inspect the original request from the sub-workflow
original_data = request.source_event.data
# Create and send a response back to the sub-workflow
response = request.create_response(my_response_data)
await ctx.send_message(response, target_id=request.executor_id)
A create_response() metódus ellenőrzi, hogy a válasz adattípusa megegyezik-e az eredeti kérés várt típusával. Ha a típusok nem egyeznek, akkor TypeError kivétel keletkezik.
Important
Amikor visszaküldi a választ, használja a target_id=request.executor_id-t a SubWorkflowResponseMessage megfelelő WorkflowExecutor példányra való irányításához.
Kérelmek propagálása külső hívóknak
Ezzel propagate_request=Trueaz al-munkafolyamatból érkező kéréseket a rendszer a standard request_info mechanizmus használatával propagálja a szülő munkafolyamat eseményfolyamára. A szülő munkafolyamat hívója ugyanúgy kezeli ezeket a kéréseket, mint bármely más, emberi beavatkozást igénylő kérést.
sub_workflow_executor = WorkflowExecutor(
workflow=inner_workflow,
id="analysis-pipeline",
propagate_request=True,
)
# Run the parent workflow and handle propagated requests
result = await parent_workflow.run(input_data)
request_info_events = result.get_request_info_events()
if request_info_events:
responses = {}
for event in request_info_events:
# Handle each request (e.g., ask a human)
responses[event.request_id] = get_human_response(event.data)
result = await parent_workflow.run(responses=responses)
Hogyan működik?
Amikor a szülő munkafolyamat átirányít egy üzenetet a WorkflowExecutorkövetkezőhöz:
- Bemeneti kézbesítés – az üzenet a belső munkafolyamat kezdő végrehajtójának továbbítja. Az üzenettípusnak meg kell egyeznie a kezdő végrehajtó várt bemeneti típusával.
- Belső végrehajtás – a belső munkafolyamat a saját superstep ciklusát futtatja a befejezésig, vagy amíg külső bemenetre nem van szüksége.
-
Kimeneti gyűjtemény – a belső munkafolyamat kimeneti eseményei a beállítás alapján
allow_direct_outputlesznek összegyűjtve és továbbítva. -
Kérelmek továbbítása – ha a belső munkafolyamat függőben lévő kérésekkel rendelkezik, a rendszer a
propagate_requestbeállítás alapján továbbítja őket (lásd : Kérések és válaszok). -
Válaszfelhalmozódás – a
WorkflowExecutorválaszok összegyűjtése és az al-munkafolyamat folytatása csak akkor történik meg, ha egy adott végrehajtásra vonatkozó összes várt válasz érkezett. - Alsóbb rétegbeli küldés – a kimenetek a szülő munkafolyamat következő végrehajtójának lesznek elküldve.
Az al-munkafolyamat a szülőtől függetlenül fenntartja saját belső állapotát. Az üzenetek csak a szülőgráf többi részéhez csatlakozó WorkflowExecutor éleken keresztül vannak irányítva – a beágyazott szinteken nem jelennek meg üzenetek.
Többszintű beágyazás
Az al-munkafolyamatok tetszőleges mélységbe ágyazhatók. Minden szint saját végrehajtási környezetet tart fenn:
# Level 1: Data preparation pipeline
data_pipeline = (
WorkflowBuilder(start_executor=fetcher)
.add_edge(fetcher, cleaner)
.build()
)
data_pipeline_executor = WorkflowExecutor(data_pipeline, id="data_pipeline")
# Level 2: Analysis pipeline (contains the data pipeline)
analysis_pipeline = (
WorkflowBuilder(start_executor=data_pipeline_executor)
.add_edge(data_pipeline_executor, analyzer)
.build()
)
analysis_pipeline_executor = WorkflowExecutor(analysis_pipeline, id="analysis_pipeline")
# Level 3: Top-level orchestration
top_workflow = (
WorkflowBuilder(start_executor=coordinator)
.add_edge(coordinator, analysis_pipeline_executor)
.add_edge(analysis_pipeline_executor, reporter)
.build()
)
Megjegyzés:
Minden beágyazási szint végrehajtási többletterhelést eredményez, mivel a belső munkafolyamat saját szupersztep-ciklust futtat. Tartsa ésszerű határok között a fészkelés mélységét a teljesítményérzékeny forgatókönyvekben.
Warning
WorkflowExecutor összes egyidejű végrehajtása ugyanazt a mögöttes munkafolyamat-példányt osztja meg. Az al-munkafolyamat végrehajtóinak állapot nélkülinek kell lenniük az egyidejű végrehajtások közötti interferencia elkerülése érdekében.
Hibakezelés
Ha egy al-munkafolyamat meghiúsul, a hiba propagálása a szülő munkafolyamatba történik. A WorkflowExecutor rendszer rögzíti a sikertelen eseményt az al-munkafolyamatból, és a szülőkörnyezetben hibaeseménysé alakítja:
async for event in parent_workflow.run(input_data, stream=True):
if event.type == "error":
print(f"Sub-workflow failed: {event.details.message}")
elif event.type == "output":
print(event.data)
Ha az al-munkafolyamat kezeletlen kivételt tapasztal, a szülő munkafolyamat hibaüzenetet kap a kivétel részleteivel, beleértve az al-munkafolyamat azonosítóját is.
Ellenőrző pontok használata
Az al-munkafolyamatok támogatják az ellenőrzőpont-készítést. Ha ellenőrzőpontot hoznak létre a szülő munkafolyamaton, a WorkflowExecutor szerializálja annak belső állapotát, beleértve a belső munkafolyamat végrehajtási állapotát és a gyorsítótárazott üzeneteket. A visszaállítás során ez az állapot deszerializálva van, lehetővé téve, hogy a szülő munkafolyamat a meglévő al-munkafolyamattal együtt folytatódjon.
from agent_framework import FileCheckpointStorage, WorkflowBuilder
checkpoint_storage = FileCheckpointStorage(storage_path="./checkpoints")
# Build the parent workflow with checkpointing
parent_workflow = (
WorkflowBuilder(
start_executor=coordinator,
checkpoint_storage=checkpoint_storage,
)
.add_edge(coordinator, inner_workflow_executor)
.add_edge(inner_workflow_executor, reviewer)
.build()
)
# Run with automatic checkpointing
async for event in parent_workflow.run("Analyze the dataset", stream=True):
if event.type == "output":
print(event.data)
# Resume from a checkpoint
checkpoints = await checkpoint_storage.list_checkpoints(workflow_name=parent_workflow.name)
async for event in parent_workflow.run(
checkpoint_id=checkpoints[-1].checkpoint_id,
checkpoint_storage=checkpoint_storage,
stream=True,
):
if event.type == "output":
print(event.data)
Sub-Workflow létrehozása
A Go nyelvben úgy hozhat létre al-munkafolyamatot, hogy létrehoz egy *workflow.Workflow elemet, majd a inproc.BindSubworkflowAsExecutor használatával a szülő munkafolyamathoz köti.
package main
import (
"context"
"fmt"
"slices"
"strings"
"github.com/microsoft/agent-framework-go/workflow"
"github.com/microsoft/agent-framework-go/workflow/inproc"
)
func buildParentWorkflow() (*workflow.Workflow, error) {
uppercase := workflow.NewExecutor("UppercaseExecutor", strings.ToUpper).Bind()
reverse := workflow.NewExecutor("ReverseExecutor", reverseString).Bind()
appendSuffix := workflow.NewExecutor("AppendSuffixExecutor", func(input string) string {
return input + " [PROCESSED]"
}).Bind()
textProcessing, err := workflow.NewBuilder(uppercase).
AddEdge(uppercase, reverse).
AddEdge(reverse, appendSuffix).
WithOutputFrom(appendSuffix).
Build()
if err != nil {
return nil, err
}
textProcessingExecutor := inproc.BindSubworkflowAsExecutor(
textProcessing,
"TextProcessingSubWorkflow",
)
prefix := workflow.NewExecutor("PrefixExecutor", func(input string) string {
return "INPUT: " + input
}).Bind()
postProcess := workflow.NewExecutor("PostProcessExecutor", func(input string) string {
return "[FINAL] " + input + " [END]"
}).Bind()
return workflow.NewBuilder(prefix).
AddEdge(prefix, textProcessingExecutor).
AddEdge(textProcessingExecutor, postProcess).
WithOutputFrom(postProcess).
Build()
}
func reverseString(input string) string {
runes := []rune(input)
slices.Reverse(runes)
return string(runes)
}
func runWorkflow(ctx context.Context, parentWorkflow *workflow.Workflow) error {
run, err := inproc.Default.RunStreaming(ctx, parentWorkflow, "hello")
if err != nil {
return err
}
defer run.Close(ctx)
for event, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
if output, ok := event.(workflow.OutputEvent); ok {
fmt.Println(output.Output)
}
}
return nil
}
A kötött gyermek munkafolyamat egyetlen végrehajtóként fut a szülő-munkafolyamat szemszögéből. Az üzenetek a kötésen keresztül lépnek be, a gyermek munkafolyamat saját belső gráfot futtat, a gyermek munkafolyamat kimenetei pedig vissza lesznek irányítva a szülőgráfba.
Bemeneti és kimeneti típusok
A kötés örökli a burkolt munkafolyamat protokollját. A szülő munkafolyamat küldhet olyan üzeneteket, amelyek futtatókörnyezeti típusai megegyeznek a gyermek munkafolyamat elfogadott bemeneti típusaival, a kötés pedig a gyermek munkafolyamat által előállított kimeneti típusokat teszi elérhetővé üzenettípusként és kimeneti típusként is.
Ez azt jelenti, hogy a szülő munkafolyamat éleinek össze kell kapcsolniuk azokat a végrehajtókat, akiknek a kimeneti típusai megegyeznek az al-munkafolyamat által elfogadott bemenetekkel, az alsóbb rétegbeli végrehajtóknak pedig a gyermek munkafolyamat által létrehozott típusokat kell kezelnie:
type TextProcessingRequest struct {
Text string
}
type TextProcessingResult struct {
Text string
}
orchestrator := workflow.NewExecutor("Orchestrator", func(ctx *workflow.Context, texts []string) error {
for _, text := range texts {
if err := ctx.SendMessage("", TextProcessingRequest{Text: text}); err != nil {
return err
}
}
return nil
}).Bind()
collector := workflow.NewExecutor("Collector", func(result TextProcessingResult) {
fmt.Println(result.Text)
}).Bind()
Kimeneti viselkedés
Amikor egy gyermek-munkafolyamat kimenetet állít elő, az almunkafolyamat-kötés ezt a kimenetet üzenetként küldi el a kötéstől a szülő munkafolyamat kapcsolódó végrehajtóinak. Ha a szülő munkafolyamat az al-munkafolyamat kötését WithOutputFromis megjelöli, ugyanaz az érték lesz kibocsátva, mint egy szülő workflow.OutputEvent , akinek ExecutorID az al-munkafolyamat kötésazonosítója.
subWorkflowExecutor := inproc.BindSubworkflowAsExecutor(textProcessing, "TextProcessingSubWorkflow")
postProcess := workflow.NewExecutor("PostProcessExecutor", func(input string) string {
return "[FINAL] " + input
}).Bind()
parentWorkflow, err := workflow.NewBuilder(subWorkflowExecutor).
AddEdge(subWorkflowExecutor, postProcess).
WithOutputFrom(subWorkflowExecutor).
WithOutputFrom(postProcess).
Build()
A gyermek munkafolyamatban kibocsátott egyéni munkafolyamat-eseményeket a rendszer a szülőesemény-adatfolyamba továbbítja. A részmunkafolyamat saját indítási és szuperlépés-életciklushoz kapcsolódó eseményei belső használatúak maradnak, így a szülőfolyam a kívülről is releváns eseményekre összpontosíthat.
Kérelmek és válaszok
Az al-munkafolyamatok támogatják a kérés- és válaszmechanizmust . Amikor egy végrehajtó az al-munkafolyamatban közzétesz egy külső kérést, az al-munkafolyamat kötése a kötésazonosító előerősítésével minősíti a kérelemport azonosítóját. A gyermekkérelem-port ApprovalPort neve például a szülő munkafolyamatban lesz ApprovalSubWorkflow.ApprovalPort .
Ha a gyermekkérelmet a szülő-munkafolyamaton keresztül szeretné felszínre helyezni, adjon hozzá egy szülőt RequestPort , amely rendelkezik a munkafolyamat-rész-kötés és a port közötti minősített azonosítóval és útvonalkérésekkel és válaszokkal:
import "reflect"
approvalPort := workflow.RequestPort{
ID: "ApprovalPort",
Request: reflect.TypeFor[string](),
Response: reflect.TypeFor[bool](),
}
approvalWorkflow, err := workflow.NewBuilder(approvalPort.Bind()).
Build()
if err != nil {
return err
}
approvalSubWorkflow := inproc.BindSubworkflowAsExecutor(
approvalWorkflow,
"ApprovalSubWorkflow",
)
qualifiedApprovalPort := workflow.RequestPort{
ID: "ApprovalSubWorkflow.ApprovalPort",
Request: approvalPort.Request,
Response: approvalPort.Response,
}
qualifiedApproval := qualifiedApprovalPort.Bind()
parentWorkflow, err := workflow.NewBuilder(approvalSubWorkflow).
AddDirectEdge(approvalSubWorkflow, qualifiedApproval, false, externalRequestOnly).
AddDirectEdge(qualifiedApproval, approvalSubWorkflow, false, externalResponseOnly).
Build()
A hívó kezeli a szülő munkafolyamat-adatfolyamból érkező kérést, és ugyanazon a lefuttatási leírón keresztül küldi vissza a választ. Az al-munkafolyamat-kötés eltávolítja a minősített névtérelőtagot, mielőtt továbbítja a választ a gyermek-munkafolyamatnak.
run, err := inproc.Default.RunStreaming(ctx, parentWorkflow, "Approve deployment?")
if err != nil {
return err
}
defer run.Close(ctx)
for event, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
switch event := event.(type) {
case workflow.RequestInfoEvent:
response, err := event.Request.CreateResponse(true)
if err != nil {
return err
}
if err := run.SendResponse(ctx, response); err != nil {
return err
}
case workflow.OutputEvent:
fmt.Println(event.Output)
}
}
Használjon predikátumokat a kérelem és a válasz élének szűkítéséhez:
func externalRequestOnly(msg any) bool {
_, ok := msg.(*workflow.ExternalRequest)
return ok
}
func externalResponseOnly(msg any) bool {
_, ok := msg.(*workflow.ExternalResponse)
return ok
}
Hogyan működik?
Amikor a szülő munkafolyamat átirányít egy üzenetet az al-munkafolyamat kötéséhez:
- Bemeneti kézbesítés – a kötés fogadja a gyermek munkafolyamat elfogadott bemeneti típusainak megfelelő üzeneteket, és a gyermek munkafolyamat kezdő végrehajtójába alakítja őket.
- Belső végrehajtás – az alárendelt munkafolyamat ugyanabban az azonos folyamaton belüli végrehajtási környezetben fut, és saját szuperlépés-ciklust tart fenn.
-
Kimenettovábbítás — a gyermek
workflow.OutputEventértékek üzenetekként kerülnek elküldésre a kötésből a lejjebb lévő szülő-végrehajtóknak, és a szülőből is kiadódnak, ha a kötés szerepel aWithOutputFromlistában. -
Kérelemtovábbítás – a gyermekkérelmek
workflow.RequestInfoEventújra kibocsáthatók minősített portazonosítókkal, és szülőkötésekenRequestPortkeresztül irányíthatók. -
Eseménytovábbítás – a rendszer hozzáadja az egyéni gyermek munkafolyamat-eseményeket a szülőstreamhez. A hibák szülő
workflow.ErrorEventértékekként jelennek meg, miközben az al-munkafolyamat azonosítója rögzítésre kerül. - Alsóbb rétegbeli küldés – az eredményként kapott üzenetek a fölérendelt munkafolyamat élein haladnak tovább.
A gyermek munkafolyamat a szülőtől elkülönítve tartja az állapotát és az üzenetek útválasztását. Az üzenetek csak a munkafolyamat-részkötéshez csatlakoztatott éleken keresztül lépik át a határt.
Többszintű beágyazás
Az al-munkafolyamatok tetszőleges mélységbe ágyazhatók. Minden egyes alárendelt munkafolyamat kötésre kerül, mielőtt hozzáadják az azt tartalmazó munkafolyamathoz:
fraudCheck, err := workflow.NewBuilder(analyzePatterns).
AddEdge(analyzePatterns, calculateRiskScore).
WithOutputFrom(calculateRiskScore).
Build()
if err != nil {
return err
}
fraudCheckExecutor := inproc.BindSubworkflowAsExecutor(fraudCheck, "FraudCheck")
payment, err := workflow.NewBuilder(validatePayment).
AddEdge(validatePayment, fraudCheckExecutor).
AddEdge(fraudCheckExecutor, chargePayment).
WithOutputFrom(chargePayment).
Build()
if err != nil {
return err
}
paymentExecutor := inproc.BindSubworkflowAsExecutor(payment, "Payment")
shippingExecutor := inproc.BindSubworkflowAsExecutor(shipping, "Shipping")
orderWorkflow, err := workflow.NewBuilder(orderReceived).
AddEdge(orderReceived, paymentExecutor).
AddEdge(paymentExecutor, shippingExecutor).
AddEdge(shippingExecutor, orderCompleted).
WithOutputFrom(orderCompleted).
Build()
Megjegyzés:
Minden egyes beágyazási szint végrehajtási többletterheléssel jár, mivel a gyermekmunkafolyamat a saját szupersztep-ciklusát futtatja. Tartsa ésszerű határok között a fészkelés mélységét a teljesítményérzékeny forgatókönyvekben.
Hibakezelés
Ha egy gyermek-munkafolyamat hibát ad vissza, az al-munkafolyamat kötése workflow.ErrorEvent formájában továbbítja azt a szülő munkafolyamathoz, és a SubWorkflowID értékét a kötésazonosítóra állítja. A szülő munkafolyamat ugyanazokat a hibákat figyelheti meg, mint a legfelső szintű munkafolyamat-hibákhoz használt eseménystream:
for event, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
switch event := event.(type) {
case workflow.ErrorEvent:
if event.SubWorkflowID != "" {
return fmt.Errorf("sub-workflow %q failed: %w", event.SubWorkflowID, event.Error)
}
return event.Error
case workflow.ExecutorFailedEvent:
return fmt.Errorf("executor %q failed: %w", event.ExecutorID, event.Error)
}
}
A gyermekesemények továbbítása során felmerülő hibák a hozzájuk csatolt al-munkafolyamat-azonosítóval a szülő workflow.ErrorEvent értékeivé is konvertálódnak.
Ellenőrző pontok használata
Az al-munkafolyamatok támogatják az ellenőrzőpont-készítést. Amikor a szülő munkafolyamat ellenőrzőpontot hoz létre, az almunkafolyamat-kötés a gyermek-munkafolyamat ellenőrzőpont-kezelőjét és a függőben lévő minősített válaszportok leképezéseit a szülő végrehajtójának állapotában tárolja. A visszaállítás során a gyermek munkafolyamat a beágyazott végrehajtási állapottal folytatható, beleértve a függőben lévő kéréseket is.
checkpointManager := checkpoint.NewInMemoryManager()
environment := inproc.Default.WithCheckpointing(checkpointManager)
var checkpoints []workflow.CheckpointInfo
run, err := environment.RunStreaming(ctx, parentWorkflow, "hello")
if err != nil {
return err
}
defer run.Close(ctx)
for event, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
if completed, ok := event.(workflow.SuperStepCompletedEvent); ok {
if completed.CompletionInfo != nil && completed.CompletionInfo.CheckpointInfo != nil {
checkpoints = append(checkpoints, *completed.CompletionInfo.CheckpointInfo)
}
}
}
if len(checkpoints) == 0 {
return fmt.Errorf("no checkpoints were created")
}
resumedRun, err := environment.ResumeStreaming(ctx, parentWorkflow, checkpoints[len(checkpoints)-1])
if err != nil {
return err
}
defer resumedRun.Close(ctx)
Ha egy ellenőrzőpont vissza lett állítva, miközben a gyermek munkafolyamat függőben lévő kéréssel rendelkezik, a visszaállított szülőfuttatás újra közzéteheti a minősített kéreleminformációs eseményt. A hívó létrehozhat egy választ abból az újraközölt kérésből, majd visszaküldheti azt a szülő futtatás kezelőjén keresztül.