使用 Temporal 作为 mcp-agent 工作流执行引擎:从部署到五种生产级工作流模式实战
2026/9/16 18:12:43 网站建设 项目流程

使用 Temporal 作为 mcp-agent 工作流执行引擎:从部署到五种生产级工作流模式实战

【免费下载链接】mcp-agentBuild effective agents using Model Context Protocol and simple workflow patterns项目地址: https://gitcode.com/GitHub_Trending/mc/mcp-agent

本篇指南围绕 mcp-agent 仓库中 examples/temporal 的完整示例集展开,系统讲解如何将 Temporal 配置为 mcp-agent 的持久化执行引擎,并依次剖析 basic、evaluator-optimizer、orchestrator、parallel、router 五种工作流模式的源码实现与运行方式。读完本文,你将掌握从启动 Temporal Server、注册 Worker,到用@app.workflow定义可暂停、可恢复、可重试的工作流,再到产出graded_report.md分级报告的完整闭环能力。

为什么选择 Temporal 作为执行引擎

mcp-agent 原生支持两种执行模式:asynciotemporal,二者的切换只需修改 mcp_agent.config.yaml 中的execution_engine配置项,一行配置即可完成引擎替换。Temporal 是微服务编排平台,为工作流提供了**持久化执行(durable execution)**能力:

  • 工作流可以长时间运行,不受进程重启影响;
  • 支持暂停、恢复、重试等运维操作,由 Temporal 平台统一保障;
  • 相同的编排能力虽然可以通过进程内(in-proc)的asyncio模式实现,但官方文档明确建议:生产环境的 mcp-agent 部署应使用工作流编排后端(即 Temporal)。

从源码角度看,TemporalExecutor是这一能力的核心实现,位于 src/mcp_agent/executor/temporal/init.py,其类注释将其职责定义为“将@workflow作为 Temporal 工作流、将@workflow_tasks作为 Temporal 活动(activities)运行的执行器”。这意味着你使用装饰器声明的业务逻辑会被翻译为 Temporal 的持久化单元,从而获得平台级的可靠性保证。

前置条件与 Temporal Server 部署

运行本示例前需要准备:

  • Python 3.10+
  • uv包管理器(示例的依赖安装与脚本执行均基于uv
  • 一个正在运行的 Temporal Server

启动 Temporal Server 最快捷的方式是使用 Temporal CLI:

temporal server start-dev

该命令会在localhost:7233启动一个开发模式服务端——这正是 mcp_agent.config.yaml 中配置的默认地址。同时可访问http://localhost:8233打开 Temporal Web UI,实时监控工作流的执行状态、历史事件与调度情况。

配置解读:execution_engine、Temporal 参数与 MCP Server

examples/temporal/mcp_agent.config.yaml 是整套示例的配置中枢,核心内容如下:

# Set the execution engine to Temporal execution_engine: "temporal" # Temporal settings temporal: host: "localhost:7233" # Default Temporal server address namespace: "default" # Default Temporal namespace task_queue: "mcp-agent" # Task queue for workflows and activities max_concurrent_activities: 10 # Maximum number of concurrent activities rpc_metadata: X-Client-Name: "mcp-agent" mcp: servers: fetch: command: "uvx" args: ["mcp-server-fetch"] description: "Fetch content at URLs from the world wide web" filesystem: command: "npx" args: ["-y", "@modelcontextprotocol/server-filesystem"] description: "Read and write files on the filesystem" openai: default_model: "gpt-4o-mini"

各参数含义与作用如下:

配置项示例值说明
execution_enginetemporal切换执行引擎的核心开关,可选asynciotemporal
temporal.hostlocalhost:7233Temporal Server 地址,与temporal server start-dev默认端口一致
temporal.namespacedefaultTemporal 命名空间,用于隔离不同业务的工作流
temporal.task_queuemcp-agent任务队列名,Worker 与客户端必须使用同一队列才能匹配
temporal.max_concurrent_activities10最大并发活动数,控制 Worker 同时执行的活动数量上限
temporal.rpc_metadataX-Client-Name: "mcp-agent"附加到 RPC 请求的自定义元数据,便于在 Temporal 侧识别客户端
logger.transports[console, file]日志输出到控制台与 JSONL 文件,path_pattern定义按时间戳命名日志文件
mcp.serversfetchfilesystem声明工作流可用的 MCP Server,fetch负责抓取网页,filesystem负责读写本地文件

值得注意的细节:filesystem服务器未在配置中预置目录参数,而是由各示例在运行时通过context.config.mcp.servers["filesystem"].args.extend([os.getcwd()])动态注入当前工作目录,例如 basic.py 与 orchestrator.py 中的写法,这样示例在不同目录下运行都能指向正确的文件系统根路径。

配置校验依赖仓库根目录下的 schema/mcp-agent.config.schema.json,配置顶部的$schema字段即指向该 JSON Schema。

四步运行:依赖、Server、Worker、示例脚本

完整运行任何一个示例都需要四个步骤:

第一步:安装依赖

uv pip install -r requirements.txt

requirements.txt 声明了四类依赖:核心框架mcp-agent(通过file://../../链接到本地仓库根目录)、anthropicopenai以及 Temporal 的 Python SDKtemporalio

第二步:启动 Temporal Server(如前述temporal server start-dev

第三步:在独立终端启动 Worker

uv run run_worker.py

Worker 会注册全部工作流并持续等待任务执行。run_worker.py 的实现极为简洁:

import workflows # noqa: F401 from main import app from mcp_agent.executor.temporal import create_temporal_worker_for_app async def main(): async with create_temporal_worker_for_app(app) as worker: await worker.run()

关键在于import workflows这一行:workflows.py 集中导入了五个示例中定义的全部工作流类与编排函数,Worker 进程只有先完成这些导入,Temporal 才能注册并识别对应的工作流类型。create_temporal_worker_for_app函数正是源码 src/mcp_agent/executor/temporal/init.py 中提供的 worker 构造入口。

第四步:在另一个终端运行示例脚本

uv run basic.py # OR uv run evaluator_optimizer.py # OR uv run orchestrator.py # OR uv run parallel.py # OR uv run router.py

五种工作流模式源码解析

所有示例共享同一个应用入口 main.py:

from mcp_agent.app import MCPApp # Create the app with Temporal as the execution engine app = MCPApp(name="temporal_workflow_example")

MCPApp会根据配置文件自动加载temporal执行引擎,示例之间仅在工作流定义上存在差异。

1. 基础工作流(basic.py):掌握装饰器三件套

basic.py 演示了 Temporal 工作流的最小骨架,也是理解后续所有示例的基石:

@app.workflow class SimpleWorkflow(Workflow[str]): @app.workflow_run async def run(self, input: str) -> WorkflowResult[str]: finder_agent = Agent( name="finder", instruction="""You are a helpful assistant.""", server_names=["fetch", "filesystem"], ) context = app.context context.config.mcp.servers["filesystem"].args.extend([os.getcwd()]) async with finder_agent: finder_llm = await finder_agent.attach_llm(OpenAIAugmentedLLM) result = await finder_llm.generate_str(message=input) return WorkflowResult(value=result) async def main(): async with app.run() as agent_app: executor: TemporalExecutor = agent_app.executor handle = await executor.start_workflow( "SimpleWorkflow", "Print the first 2 paragraphs of https://modelcontextprotocol.io/introduction", ) a = await handle.result() print(a)

核心要素拆解:

  • @app.workflow:声明该类为 Temporal 工作流,泛型参数Workflow[str]表示输入类型;
  • @app.workflow_run:标记run方法为工作流执行体,接收输入并返回WorkflowResult[str]包装的结果;
  • 工作流内部创建 Agentfinder代理绑定fetchfilesystem两个 MCP Server,在async with上下文中完成 LLM 的挂载(attach_llm(OpenAIAugmentedLLM))与生成调用(generate_str);
  • 执行与等待agent_app.executor暴露TemporalExecutor,调用start_workflow("SimpleWorkflow", input)返回句柄,await handle.result()阻塞等待工作流在 Temporal 侧完成并取回结果。start_workflow的签名定义见 src/mcp_agent/executor/temporal/init.py。

2. Evaluator-Optimizer 工作流:让反馈循环驱动内容迭代

evaluator_optimizer.py 模拟求职场景:根据职位描述、候选人信息与公司资料生成求职信,再由评估者持续评判、迭代打磨,直到达到质量门槛。

optimizer = Agent( name="optimizer", instruction="""You are a career coach specializing in cover letter writing. ...""", server_names=["fetch"], ) evaluator = Agent( name="evaluator", instruction="""Evaluate the following response based on the criteria below: 1. Clarity: ... 2. Specificity: ... 3. Relevance: ... 4. Tone and Style: ... 5. Persuasiveness: ... 6. Grammar and Mechanics: ... 7. Feedback Alignment: ... For each criterion: - Provide a rating (EXCELLENT, GOOD, FAIR, or POOR). - Offer specific feedback or suggestions for improvement. Summarize your evaluation as a structured response with: - Overall quality rating. - Specific feedback and areas for improvement.""", ) evaluator_optimizer = EvaluatorOptimizerLLM( optimizer=optimizer, evaluator=evaluator, llm_factory=OpenAIAugmentedLLM, min_rating=QualityRating.EXCELLENT, context=app.context, ) result = await evaluator_optimizer.generate_str( message=input, request_params=RequestParams(model="gpt-4o"), ) return WorkflowResult(value=result)

该模式的关键设计:评估者指令中内置了七大评判维度(清晰度、具体性、相关性、语气风格、说服力、语法、反馈对齐度),并要求给出EXCELLENT / GOOD / FAIR / POOR分级评价;EvaluatorOptimizerLLM将二者组装为反馈闭环,min_rating=QualityRating.EXCELLENT设定迭代停止条件——只有达到最高评级才结束循环。通过RequestParams(model="gpt-4o")可为单次生成覆盖默认模型。

3. Orchestrator 工作流:动态编排多代理协作

orchestrator.py 展示了更复杂的多代理编排:它不直接使用@app.workflow+@app.workflow_run,而是改用@app.async_tool装饰器,将整个编排逻辑封装为一个异步工具,工作流由框架在后台自动创建:

@app.async_tool(name="OrchestratorWorkflow") async def run_orchestrator(input: str, app_ctx: Optional[AppContext] = None) -> str: ... orchestrator = Orchestrator( llm_factory=OpenAIAugmentedLLM, available_agents=[ finder_agent, writer_agent, proofreader, fact_checker, style_enforcer, ], plan_type="full", # 每一步都由编排器动态规划 context=context, ) return await orchestrator.generate_str( message=input, request_params=RequestParams(model="gpt-4o", max_iterations=100), )

五个 Agent 分工明确:finder负责检索文件与网页并返回最近匹配项的 URI 与内容;writer将结果写入磁盘;proofreaderfact_checkerstyle_enforcer分别从语法、事实一致性、风格规范三个角度审阅;plan_type="full"意味着编排器会在每一步动态决定下一步动作。主流程给它下达的任务是:读取 short_story.md 中的学生短篇小说,参考 APA 风格指南生成涵盖校对、事实逻辑与风格遵循的分级报告,并写入 graded_report.md。

4. Parallel 工作流:Fan-out/Fan-in 并行处理

parallel.py 演示扇出/扇入(fan-out/fan-in)模式:将同一篇短篇小说同时分发给校对、事实核查、风格强化三个专家代理并行处理,再由grader汇总为结构化报告。

parallel = ParallelLLM( fan_in_agent=grader, fan_out_agents=[proofreader, fact_checker, style_enforcer], llm_factory=OpenAIAugmentedLLM, context=app.context, ) result = await parallel.generate_str( message=f"Student short story submission: {input}", )

该示例还附带了token 用量统计的完整实践:

  • 工作流内部通过parallel.get_token_node()获取 token 树、用量与成本,写入WorkflowResultmetadata字段返回;
  • 主进程除打印返回结果外,还尝试通过handle.query("token_tree")handle.query("token_summary")运行中的工作流发起 Temporal 查询,实时获取远程工作流内的 token 统计——这是 Temporal 查询(Query)机制在 mcp-agent 中的典型用法;
  • 代码注释同时提醒:Temporal 模式下客户端进程内的TokenCounter可能为 0,应优先依赖工作流侧查询的指标。

5. Router 工作流:面向 Agent、函数与 Server 的智能路由

router.py 展示了 LLM 驱动的路由决策能力,且覆盖四种路由场景:

llm = OpenAIAugmentedLLM(name="openai_router", instruction="You are a router") router = LLMRouter( llm_factory=lambda _agent: llm, agents=[finder_agent, writer_agent, reasoning_agent], functions=[print_to_console, print_hello_world], context=app.context, ) # 路由到 Agent:读取 mcp_agent.config.yaml 内容 results = await router.route_to_agent(request="...", top_k=1) # 路由到函数:打印输入到控制台 results = await anthropic_router.route_to_function(request="...", top_k=2) function_to_call = results[0].result function_to_call("Hello, world!") # 仅凭 Server 名称推断路由目标 results = await anthropic_router.route_to_server( request="Print the first two paragraphs of https://modelcontextprotocol.io/introduction", top_k=1, ) # 跨所有类别统一路由(servers + agents + callables) results = await anthropic_router.route(request="...", top_k=3)

四种路由 API 的区别与适用场景:

路由方法路由目标范围典型场景
route_to_agent仅限 Agent 列表把请求分派给最合适的代理执行
route_to_function仅限普通 Python 函数把请求映射到本地可调用对象
route_to_server仅限 MCP Server根据 Server 名称/描述推断目标工具来源
route所有类别混合跨 Agent、函数、Server 的统一全局路由

示例同时给出两个路由器实现:LLMRouter(基于 OpenAI 系 LLM,可通过llm_factory注入任意 LLM)与AnthropicLLMRouter(预配置 Anthropic LLM 并绑定fetch/filesystem两个 Server)。top_k控制返回候选数量,且示例注释指出路由会依据请求内容做语义匹配——例如请求“打印到控制台”时即便top_k=2也只返回print_to_console而不会误返回print_hello_world

项目结构总览

examples/temporal/ ├── main.py # 核心应用配置:MCPApp 与执行引擎 ├── run_worker.py # Worker 启动脚本(create_temporal_worker_for_app) ├── workflows.py # 集中导入全部工作流,供 Worker 注册 ├── basic.py # 基础工作流示例 ├── evaluator_optimizer.py # 评估-优化迭代示例 ├── orchestrator.py # 多代理动态编排示例 ├── parallel.py # 并行扇出/扇入示例 ├── router.py # 智能路由示例 ├── interactive.py # 带交互的工作流示例 ├── short_story.md # 示例使用的学生短篇小说样本 ├── graded_report.md # orchestrator/parallel 工作流的输出报告 ├── mcp_agent.config.yaml # 引擎与 Temporal 配置 └── requirements.txt # 依赖清单

工作流生命周期总结

将上述示例串成一条完整链路,可以清晰看到 mcp-agent + Temporal 的协作模型:

  1. 定义:在任意示例脚本中用@app.workflow/@app.workflow_run(或@app.async_tool)声明工作流;
  2. 注册:Worker 进程通过import workflows加载所有工作流定义,并由create_temporal_worker_for_app(app)注册到 Temporal;
  3. 启动:客户端进程在app.run()上下文中取得TemporalExecutor,调用executor.start_workflow(name, input)向任务队列mcp-agent提交执行;
  4. 执行与恢复:Temporal 按事件驱动执行工作流,活动失败自动重试,进程重启后可从历史事件恢复,这正是“持久化执行”在生产环境的核心价值;
  5. 取回结果await handle.result()等待完成,必要时可通过handle.query(...)对运行中的工作流发起实时查询(如 token 统计)。

如果想验证更复杂的交互式场景,还可运行 interactive.py 了解如何在 Temporal 工作流中插入人工输入环节。生产环境若要深入掌控 executor 的并发、重试与查询行为,可直接研读 src/mcp_agent/executor/temporal 下的TemporalExecutor实现与其客户端拦截器(interceptor)。

【免费下载链接】mcp-agentBuild effective agents using Model Context Protocol and simple workflow patterns项目地址: https://gitcode.com/GitHub_Trending/mc/mcp-agent

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询