子工作流

子工作流是作为父工作流中的执行程序运行的完整工作流。 这使你能够从较小的可重用工作流构建基块(每个系统都有自己的独立执行上下文、状态管理和消息路由)组成复杂的系统。

Overview

在以下情况下需要时,子工作流非常有用:

  • 分解复杂性 - 将大型工作流分解为较小的独立可测试单元。
  • 重用工作流逻辑 - 将同一子工作流嵌入多个父工作流。
  • 隔离状态 - 使每个子工作流的内部状态与父工作流保持隔离。
  • 控制数据流 - 消息仅通过其边缘进入和离开子工作流,且不跨级别广播。

将子工作流添加到父工作流时,其行为类似于任何其他执行程序:它接收输入消息,运行其内部图以完成,并为下游执行程序生成输出消息。

创建子工作流

在 C# 中,可通过两种方式编写子工作流:

  • 直接绑定 - 用于 BindAsExecutor() 将工作流直接嵌入父工作流中的执行程序。 这会保留子工作流的本机输入/输出类型。
  • 代理包装 - 用于 AsAIAgent() 将工作流转换为代理,然后将代理添加到父工作流。 当父工作流使用基于代理的执行程序时,这非常有用。

使用 BindAsExecutor 进行直接绑定

BindAsExecutor()扩展方法将工作流转换为ExecutorBinding,可直接添加到父工作流中:

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();

使用 BindAsExecutor时,子工作流的类型化输入和输出类型将保留 — 父工作流根据子工作流预期的实际类型路由消息并生成。

使用 AsAIAgent 包装代理

当父工作流使用基于代理的执行程序时,使用 AsAIAgent() 将内部工作流转换为代理。 WorkflowBuilder 会自动将代理包装在执行器中。

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();

从父工作流的角度来看,内部工作流作为单个步骤运行。 协调器将消息发送到分析管道,该管道在内部运行 specialist1 → specialist2,然后将结果转发给审阅者。

Tip

在使用类型化执行程序时使用BindAsExecutor(),在处理基于代理的工作流时使用AsAIAgent()。 有关配置工作流到代理转换的详细信息,请参阅 工作流作为代理

输入和输出类型

当工作流用作子工作流时,它将保留其内部执行程序的类型协定。

使用 BindAsExecutor时,子工作流执行程序接受与内部工作流的启动执行程序相同的输入类型,并发送内部工作流生成的相同输出类型。 父工作流的边缘必须连接其输出类型与子工作流的预期输入类型匹配的执行程序,而子工作流的输出类型必须与下游执行程序的预期输入匹配。

使用AsAIAgent时,子工作流包装为代理,并遵循代理执行器输入/输出协定 (string, , ChatMessageIEnumerable<ChatMessage>)。

输出行为

默认情况下,当子工作流生成输出时(通过 YieldOutputAsync),这些输出将作为消息转发到父工作流中连接的执行程序。 这使下游执行程序能够处理子工作流结果。

ExecutorOptions 控制此行为:

选项 默认 Description
AutoSendMessageHandlerResultObject true 将子工作流输出作为消息转发到父图中连接的执行程序。
AutoYieldOutputHandlerResultObject false 将子工作流输出直接输出到父工作流的输出事件流。

启用后 AutoYieldOutputHandlerResultObject ,子工作流输出将绕过父级的内部路由,并直接传递到父工作流的调用方。

var options = new ExecutorOptions
{
    AutoYieldOutputHandlerResultObject = true,
};

ExecutorBinding subWorkflowExecutor = innerWorkflow.BindAsExecutor("SubWorkflow", options);

请求和响应

子工作流完全支持 请求和响应 机制。 当子工作流内的执行程序发送请求(例如,请求人工输入)时,WorkflowHostExecutor会将具有RequestInfoEvent 的父工作流转发到父工作流 , 子工作流执行程序的 ID 将追加到端口 ID(例如, SubWorkflow.GuessNumber) 。

此限定可确保父工作流收到响应时,它可以将响应路由回正确的子工作流实例。 父工作流使用与任何其他请求相同的响应机制处理子工作流请求:

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;
    }
}

注释

从父工作流调用方的角度来看,来自顶级执行程序的请求与子工作流的请求之间没有区别。 框架以透明方式处理路由。

工作原理

父工作流将消息路由到子工作流执行程序时:

  1. 输入传递 - 消息转发到内部工作流的启动执行程序。 使用 BindAsExecutor时,消息类型必须与启动执行程序的预期类型匹配。 使用 AsAIAgent时,消息将规范化为 ChatMessage 格式。
  2. 内部执行 - 内部工作流运行其自己的超级步骤循环。
  3. 输出集合 - 收集内部工作流的输出事件。 使用 BindAsExecutor时,输出会保留其原始类型。 使用 AsAIAgent时,输出将转换为代理响应消息。
  4. 请求转发 - 如果内部工作流有挂起的请求,则会将其转发到父工作流进行处理(请参阅 请求和响应)。
  5. 下游调度 - 生成的消息将发送到父工作流中的下一个执行程序。

由于内部工作流维护自己的执行上下文,因此其状态独立于父工作流。

Tip

有关配置工作流到代理转换的详细信息,包括流式处理行为和异常处理,请参阅 工作流作为代理

多级嵌套

子工作流可以嵌套到任意深度。 每个级别都维护其自己的执行上下文:

// 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();

注释

每个嵌套级别都会增加执行开销,因为内部工作流运行自己的超级步骤循环。 将嵌套深度保持在合理范围,以适应对性能敏感的场景。

错误处理

当子工作流失败时,错误将作为一个 SubworkflowErrorEvent传播到父工作流。 父工作流可以通过其事件流观察这些错误:

await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
    if (evt is SubworkflowErrorEvent subError)
    {
        Console.WriteLine($"Sub-workflow '{subError.ExecutorId}' failed: {subError.Data}");
    }
}

如果子工作流遇到未经处理的异常,则父工作流的执行将继续执行,但子工作流执行程序停止处理其他消息。

检查点

在父工作流上执行检查点时,子工作流代理的会话状态将序列化为父执行程序的检查点数据的一部分。 在还原时,会话状态将反序列化,使父工作流可以恢复子工作流的状态不变。

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);

创建子工作流

在 Python 中,通过包装一WorkflowWorkflowExecutor个子工作流并将其添加到父工作流来创建子工作流。

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()
)

从父工作流的角度来看,内部工作流作为单个步骤运行。 协调器将消息发送到分析管道,该管道在内部运行 specialist1 → specialist2,然后将结果转发给审阅者。

WorkflowExecutor 参数

参数 类型 默认 Description
workflow Workflow 要包装为执行程序的工作流实例。
id str 此执行程序的唯一标识符。
allow_direct_output bool False True 子工作流的输出直接被发送到父工作流的事件流中,而不是作为消息发送给连接的执行器时。
propagate_request bool False True,从子工作流发出的请求作为常规请求信息事件传播到父工作流的事件流。 当 False 请求被包装在 SubWorkflowRequestMessage 中,以便供父执行程序拦截。

包装子工作流

在将 Workflow 实例添加到父工作流之前,先将其显式包装在 WorkflowExecutor 中。 代理可以直接传递到 WorkflowBuilder,但原始 Workflow 实例需要此包装器。

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()
)

显式包装可用于:

  • 为多个边缘中的引用分配特定的执行器ID。
  • 在整个关系图中重复使用同一 WorkflowExecutor 实例。
# 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()
)

输入和输出类型

WorkflowExecutor 继承它的类型签名自所包装的工作流。

  • 输入类型 与包装工作流的启动执行程序输入类型匹配(以及 SubWorkflowResponseMessage 用于处理对转发请求的响应)。
  • 输出类型 与包装工作流的输出类型匹配。 如果子工作流中的任何执行程序支持请求响应, SubWorkflowRequestMessage 也作为输出类型包含在内。

这意味着父工作流的边缘必须连接其输出类型与子工作流的预期输入类型匹配的执行程序。 同样,下游执行程序必须接受子工作流生成的类型:

# 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}")

输出行为

默认情况下(allow_direct_output=False),当子工作流通过yield_output生成输出时,这些输出将作为消息转发到父工作流中的连接执行程序。send_message 这使下游执行程序能够将子工作流结果作为父图的一部分进行处理。

allow_direct_output=True 时,子工作流的输出直接传送到父工作流的事件流。 子工作流的输出将成为父工作流的输出,绕过父级的内部执行程序路由:

# 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)

子工作流的中间输出

"intermediate" 由子工作流内部产生的事件会自动通过父工作流的事件流逐级上传。 它们归属于 WorkflowExecutor 自身的 id(而不是最初发出它们的内部执行器),这保留了封装性。 关键在于,无论父项如何在其自身的 "intermediate" 列表中指定 WorkflowExecutor,这些事件output_from

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}")

请求和响应

子工作流完全支持 请求和响应 机制。 当子工作流内的执行程序调用ctx.request_info()时,WorkflowExecutor会截获请求,并根据propagate_request设置进行处理。

拦截父工作流中的请求(默认值)

使用 propagate_request=False (默认设置),子工作流的请求会被包装在 SubWorkflowRequestMessage 中,并且发送到位于父工作流中的已连接执行程序。 这样,父执行程序就可以在本地处理请求:

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)

该方法 create_response() 验证响应数据类型是否与原始请求的预期类型匹配。 如果类型不匹配,则会引发TypeError

Important

在发送回响应时,使用target_id=request.executor_idSubWorkflowResponseMessage路由到正确的WorkflowExecutor实例。

向外部调用方传播请求

使用 propagate_request=True时,子工作流发出的请求将使用标准 request_info 机制传播到父工作流的事件流。 父工作流的调用方以与任何其他人循环请求相同的方式处理这些请求:

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)

工作原理

父工作流将消息路由到 WorkflowExecutor时:

  1. 输入传递 - 消息转发到内部工作流的启动执行程序。 消息类型必须与启动执行程序的预期输入类型匹配。
  2. 内部执行 — 内部工作流运行自己的超步循环以完成,或直到需要外部输入为止。
  3. 输出集合 - 根据 allow_direct_output 设置收集和转发内部工作流的输出事件。
  4. 请求转发 - 如果内部工作流具有挂起的请求,则会根据 propagate_request 设置转发它们(请参阅 请求和响应)。
  5. 响应累积 - 仅当收到给定执行的所有预期响应时,才会 WorkflowExecutor 收集响应并恢复子工作流。
  6. 下游调度 - 输出将发送到父工作流中的下一个执行程序。

子工作流独立于父工作流维护自己的内部状态。 消息仅通过连接到 WorkflowExecutor 父图的其余部分的边缘进行路由 - 没有跨嵌套级别的消息广播。

多级嵌套

子工作流可以嵌套到任意深度。 每个级别都维护其自己的执行上下文:

# 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()
)

注释

每个嵌套级别都会增加执行开销,因为内部工作流运行自己的超级步骤循环。 将嵌套深度保持在合理范围,以适应对性能敏感的场景。

Warning

所有并发执行共享同一的 WorkflowExecutor 基础工作流实例。 子工作流中的执行程序应是无状态的,以避免并发执行之间的干扰。

错误处理

当子工作流失败时,错误将传播到父工作流。 WorkflowExecutor 捕捉子工作流中的失败事件,并将其转换为父上下文中的错误事件。

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)

如果子工作流遇到未经处理的异常,则父工作流会收到异常详细信息的错误事件,包括子工作流的 ID。

检查点

子工作流支持检查点。 当检查点被保存到父工作流时,WorkflowExecutor 会将其内部状态序列化,包括内部工作流的执行进度以及任何缓存的消息。 在还原时,将反序列化此状态,使父工作流能够随着子工作流保持不变而恢复。

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)

创建子工作流

在 Go 中,您可以通过构建一个 *workflow.Workflow 并使用 inproc.BindSubworkflowAsExecutor 将其绑定到父工作流中来创建子工作流。

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
}

绑定子工作流从父工作流的角度作为一个执行程序运行。 消息通过绑定进入,子工作流运行其自身的内部图,子工作流的输出被路由回父图。

输入和输出类型

该绑定继承了所封装工作流的协议。 父工作流可以发送运行时类型与子工作流接受的输入类型匹配的消息,绑定会将子工作流生成的输出类型公开为消息类型和输出类型。

这意味着父工作流的边缘必须连接其输出类型与子工作流接受的输入匹配的执行程序,下游执行程序必须处理子工作流生成的类型:

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()

输出行为

当子工作流生成输出时,子工作流绑定会将该输出作为消息从绑定发送到父工作流中连接的执行程序。 如果父工作流还用 WithOutputFrom 标记了子工作流绑定,则会发出相同的值,作为一个父级 workflow.OutputEvent,其 ExecutorID 为子工作流绑定 ID。

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()

子工作流内发出的自定义工作流事件将转发到父事件流。 子工作流自身的启动和超级步生命周期事件仅限于内部,从而使父级事件流专注于对外部真正有意义的事件。

请求和响应

子工作流支持 请求和响应 机制。 当子工作流内的执行程序发布外部请求时,子工作流绑定通过附加绑定 ID 来限定请求端口 ID。 例如,名为 ApprovalPort 的子请求端口在父工作流中会变为 ApprovalSubWorkflow.ApprovalPort

若要通过父工作流暴露子请求,请添加一个带有限定 ID 的父级 RequestPort,并在子工作流绑定与该端口之间路由请求和响应:

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()

调用方处理来自父工作流流的请求,并通过同一运行句柄发送响应。 子工作流绑定在将响应传送到子工作流之前删除限定的前缀。

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)
    }
}

使用谓词使请求和响应边缘保持窄:

func externalRequestOnly(msg any) bool {
    _, ok := msg.(*workflow.ExternalRequest)
    return ok
}

func externalResponseOnly(msg any) bool {
    _, ok := msg.(*workflow.ExternalResponse)
    return ok
}

工作原理

父工作流将消息路由到子工作流绑定时:

  1. 输入传递 — 绑定接受与子工作流接受的输入类型匹配的消息,并将消息排入子工作流的启动执行程序。
  2. 内部执行 - 子工作流在同一进程内执行环境中运行,并维护其自己的超步循环。
  3. 输出转发 — 子级 workflow.OutputEvent 值会作为消息从该绑定发送到下游父执行器;如果该绑定列在 WithOutputFrom 中,这些值也会由父级产出。
  4. 请求转发 - 子 workflow.RequestInfoEvent 请求使用限定的端口 ID 重新发出,可以通过父 RequestPort 绑定进行路由。
  5. 事件转发 - 自定义子工作流事件将添加到父流。 错误会显示为父级 workflow.ErrorEvent 值,并记录子工作流 ID。
  6. 下游调度 - 生成的消息继续通过父工作流边缘。

子工作流将其状态和消息路由与父工作流分开。 消息仅通过连接到子工作流绑定的边缘越过边界。

多级嵌套

子工作流可以嵌套到任意深度。 每个子工作流在被添加到包含它的工作流之前,都会先完成绑定:

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()

注释

每个嵌套层级都会增加执行开销,因为子工作流会运行其自身的 superstep 循环。 将嵌套深度保持在合理范围,以适应对性能敏感的场景。

错误处理

当子工作流发出错误时,子工作流绑定会将其作为 workflow.ErrorEvent 转发给父工作流,并将 SubWorkflowID 设置为该绑定的 ID。 父工作流可以通过用于顶级工作流错误的相同事件流来观察这些错误:

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)
    }
}

转发子事件时产生的错误也会被转换为附加了子工作流 ID 的父级 workflow.ErrorEvent 值。

检查点

子工作流支持检查点。 当父工作流进行检查点保存时,子工作流绑定会将子工作流的检查点管理器以及任何待处理的限定响应端口映射存储在父执行器状态中。 在还原时,子工作流可以恢复运行,并保留其完整的嵌套执行状态,包括挂起的请求。

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)

如果在子工作流存在待处理请求时恢复检查点,则恢复后的父工作流运行会重新发布符合条件的请求信息事件。 调用方可以从该重新发布的请求创建响应,并通过父运行句柄发送回该响应。

后续步骤