三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

工具执行引擎核心设计:DAG并行、中间件与HITL重运行

工具执行引擎核心设计:DAG并行、中间件与HITL重运行

1. 从一个“简单”需求说起:为什么我们需要工具执行引擎?

如果你做过一些自动化脚本,或者写过一些需要调用外部工具(比如调用一个API、执行一个系统命令、处理一个文件)的程序,你可能会觉得这很简单:不就是按顺序写几个函数调用吗?比如,先调用工具A,拿到结果,再传给工具B,最后汇总输出。在单个任务、逻辑线性的场景下,这确实没问题。

但现实项目往往复杂得多。想象一下这个场景:你正在构建一个智能客服的对话系统。用户问:“帮我查一下明天北京的天气,然后根据天气推荐一个室内或室外的活动,最后把活动信息和天气一起总结成邮件草稿。”

拆解一下,这个任务至少涉及三个工具:

  1. 天气查询工具:调用一个天气API,获取明天北京的天气数据(温度、降水概率、风力等)。
  2. 活动推荐工具:根据天气数据(比如下雨就推荐室内博物馆,晴天就推荐户外公园),调用另一个知识库或推荐API。
  3. 邮件生成工具:将前两步的结果(天气+活动)作为输入,按照模板生成一封结构化的邮件草稿。

最直观的写法是串行:天气结果 = 调用天气工具(北京)->活动结果 = 调用推荐工具(天气结果)->邮件草稿 = 调用邮件工具(天气结果, 活动结果)

这看起来清晰,但问题马上来了:

  • 效率低下:如果“活动推荐”不依赖“天气数据”的全部细节,而只依赖一个“天气类型”(晴/雨),那么“活动推荐”工具是否可以和“天气查询”中获取“天气类型”的那部分逻辑并行?串行执行浪费了等待时间。
  • 错误处理僵化:如果“天气查询”API调用失败了,整个流程是直接报错退出,还是尝试使用缓存的历史数据?或者跳过天气,直接给一个通用的活动推荐?串行逻辑很难优雅地处理这种“部分失败”和“降级策略”。
  • 逻辑穿插困难:如果需要在每个工具调用前后都统一打印日志、计算耗时、检查权限,你需要在每个工具调用的地方重复写这些代码。更复杂的是,如果生成邮件草稿后,还需要调用一个“内容安全审核”工具,审核不通过则要触发“人工修改流程”,这个“人工介入”的环节如何自然地嵌入到自动化的流程中?

你会发现,当工具数量增多、依赖关系变复杂、需要加入统一管控(日志、鉴权)或灵活流程(人工干预、条件分支)时,简单的过程式代码会迅速变得难以维护和扩展。这时,我们就需要一个工具执行引擎来管理这些工具的调度、依赖、上下文传递和生命周期。ToolsNode正是为了解决这类问题而生的一个设计范式或框架的核心抽象。今天,我们就来彻底拆解它,聚焦于其三大核心能力:并行执行、中间件洋葱模型,以及HITL Rerun(人在回路重运行)

2. 核心模型:ToolsNode 是什么?不是什么?

在深入细节前,我们必须先统一认知。ToolsNode不是一个特指某个开源库(比如langchainTool),而是一种设计模式架构单元的抽象。你可以把它理解为一个有输入、有输出、有执行逻辑,且能被某个引擎统一调度和管理的功能单元

  • 它是什么?

    1. 功能封装单元:一个ToolsNode封装了一个具体的、可执行的操作。比如“调用天气API”、“查询数据库”、“发送HTTP请求”、“执行Python函数”。
    2. 数据流节点:它定义了明确的输入参数和输出结果。输入来自上游节点的输出或初始参数,输出会传递给下游节点作为输入。
    3. 执行调度单元:执行引擎(Orchestrator)可以并发地、按依赖顺序地触发一个或多个ToolsNode的执行。
  • 它不是什么?

    1. 它不是简单的函数。函数缺乏对依赖、生命周期、统一拦截(中间件)的内置支持。
    2. 它不是微服务。它更轻量,通常是进程内组件,通信开销低,专注于逻辑编排而非服务治理。
    3. 它不是工作流引擎中的“人工任务”。虽然它可以包含人工交互,但其核心是自动化工具。

一个典型的ToolsNode接口可能长这样(以伪代码示意):

class ToolsNode: name: str # 节点唯一标识 description: str # 节点功能描述 input_schema: dict # 输入参数JSON Schema output_schema: dict # 输出结果JSON Schema async def execute(self, input_data: dict, context: ExecutionContext) -> dict: # 核心执行逻辑 # 可以使用 context 获取共享数据、调用其他服务等 pass def get_dependencies(self) -> List[str]: # 返回本节点所依赖的其他节点名称列表 # 用于构建执行图 pass

执行引擎会根据get_dependencies()返回的依赖关系,构建一个有向无环图,然后决定哪些节点可以并行执行,哪些必须等待前置节点完成。

3. 能力一:基于DAG的智能并行执行

并行是提升复杂工具链效率的关键。但并行不是乱并行,必须尊重工具间的数据依赖ToolsNode通过依赖声明,让执行引擎能够自动推导出最优的并行方案。

3.1 依赖声明与执行图构建

每个ToolsNode都需要明确声明它依赖哪些其他节点的输出。假设我们有四个节点:

  • A(用户输入解析):无依赖。
  • B(天气查询):依赖 A 输出的location字段。
  • C(新闻检索):依赖 A 输出的topic字段。
  • D(报告生成):依赖 B 输出的weather_data和 C 输出的news_summary

它们的依赖关系可以表示为:

A -> B -> D A -> C -> D

即,B和C都依赖A,D依赖B和C。执行引擎会将其解析为一个DAG(有向无环图)。

3.2 并行策略与调度

一个智能的执行引擎会这样调度:

  1. 第一轮:发现只有 A 没有前置依赖,立即执行 A。
  2. 第二轮:A 执行完成后,B 和 C 的所有依赖(A)都已就绪。引擎会同时(并行)执行 B 和 C,因为它们之间没有依赖关系。
  3. 第三轮:B 和 C 都完成后,D 的依赖全部就绪,执行 D。

这样,总耗时从串行的A+B+C+D缩短为A + max(B, C) + D。如果 B 和 C 耗时都是2秒,串行需要A+2+2+D,并行则只需要A+2+D,节省了2秒。

实操心得:依赖声明的粒度依赖声明并非越细越好。例如,节点B声明它依赖节点A的整个输出字典。如果A的输出很大,但B只关心其中的一个小字段,这会造成不必要的数据传递和耦合。更优的做法是,让依赖声明支持路径映射,例如B.depends_on(A, output_map={"location": "A.output.user_query.location"})。这样,B的输入接口更清晰,且引擎可能有机会做更细粒度的优化。不过,这增加了复杂性,需要根据工具链的稳定性和性能要求进行权衡。

3.3 错误处理与并行安全

并行执行引入了新的复杂性:错误处理。如果并行执行的 B 和 C 中有一个失败了(比如 C 调用的新闻服务超时),D 节点应该怎么办?

  1. 快速失败:整个流程立即终止,返回错误。适用于强依赖所有前置结果的场景。
  2. 降级处理:允许节点声明某些依赖是“可选的”。如果 C 失败,D 节点可以接收到一个标记为失败或为空的结果,并在其内部逻辑中决定是否使用默认值或跳过部分功能继续执行。这需要节点逻辑有更强的鲁棒性。
  3. 重试与备用:引擎可以配置重试策略。对于 C 失败,可以重试几次,或者启用一个备用的“缓存新闻查询”节点 C‘。

ToolsNode的设计中,通常由执行引擎来提供统一的错误处理策略配置,而ToolsNode自身的execute方法应抛出结构化的异常,以便引擎捕获和决策。

4. 能力二:中间件洋葱模型——统一的横切面管控

如果你在每个ToolsNodeexecute方法里都写一遍日志、性能监控、权限校验、输入校验的代码,那将是一场维护灾难。这就是“横切关注点”问题。中间件洋葱模型是解决这个问题的经典模式。

4.1 洋葱模型如何工作?

想象一下,一个ToolsNode的执行就像一颗洋葱的核心。中间件就是一层层的洋葱皮。执行引擎在调用节点的execute方法前,会先按顺序经过一系列中间件;调用结束后,结果又会以相反的顺序再次经过这些中间件。

执行顺序: [ 中间件1前置逻辑 ] -> [ 中间件2前置逻辑 ] -> [ ToolsNode.execute() ] -> [ 中间件2后置逻辑 ] -> [ 中间件1后置逻辑 ]

这就形成了一个“洋葱”式的调用链。每个中间件都有机会在工具执行插入逻辑。

4.2 常见中间件场景与实现

让我们看几个具体的中间件例子,它们能极大提升工具链的可观测性和可控性。

1. 日志与指标中间件

class LoggingMiddleware: async def __call__(self, node: ToolsNode, input_data: dict, context: ExecutionContext, next_callable): start_time = time.time() node_name = node.name logger.info(f"开始执行节点: {node_name}, 输入: {input_data}") try: # 调用下一个中间件或最终的 node.execute result = await next_callable(node, input_data, context) duration = time.time() - start_time logger.info(f"节点执行成功: {node_name}, 耗时: {duration:.2f}s, 输出: {result}") # 可以上报指标,如 prometheus.gauge('node_duration_seconds').set(duration) return result except Exception as e: duration = time.time() - start_time logger.error(f"节点执行失败: {node_name}, 耗时: {duration:.2f}s, 错误: {e}") raise

这个中间件自动为每个节点的执行记录了开始/结束时间、输入输出和异常,无需修改任何节点代码。

2. 输入验证与转换中间件节点可能期望输入是特定格式。中间件可以根据节点的input_schema(如JSON Schema)在调用前验证输入数据,甚至进行类型转换(比如把字符串数字转成整数),确保节点核心逻辑收到的数据是干净的。

3. 权限与配额检查中间件在工具链执行前,检查当前上下文(用户、API Key)是否有权限执行该节点,或者是否超过了调用频率限制。如果无权限或超限,直接在中间件层拒绝,无需进入实际执行。

4. 缓存中间件这是性能优化的利器。中间件可以根据节点名称和输入参数的哈希值(hash(node.name + str(sorted(input_data.items()))))作为缓存键。在执行前先查缓存,命中则直接返回缓存结果,跳过节点执行;未命中则执行节点,并将结果写入缓存。

注意:缓存中间件要慎用。必须确保节点的执行是幂等的(相同输入总是产生相同输出),且缓存失效策略要合理。对于查询类、计算类工具非常适合,但对于发送邮件、写入数据库等有副作用的操作则绝对不能缓存。

实操心得:中间件的执行顺序至关重要中间件的注册顺序就是洋葱皮的包裹顺序。例如,你应该先注册权限校验中间件,再注册日志中间件。这样,如果权限校验失败,请求根本不会进入后续中间件和核心节点,但权限校验的失败日志仍然会被最外层的日志中间件记录。错误的顺序可能导致安全问题(如先缓存后鉴权)或日志缺失。

5. 能力三:HITL Rerun(人在回路重运行)——当自动化遇到瓶颈

HITL (Human-In-The-Loop) 是AI产品中常见的概念,指在自动化流程中引入人工干预点。ToolsNode框架中的HITL Rerun特指一种能力:当某个自动化节点执行失败或结果不确定时,系统能暂停流程,将问题和上下文提交给人工处理,待人工提供结果或修正后,系统能从该节点重新运行后续流程,而不是从头开始。

5.1 为什么需要 HITL Rerun?

考虑一个“自动审核用户生成内容”的流程:

  1. 节点A:内容提取(成功)
  2. 节点B:敏感词过滤(成功,标记出疑似敏感词)
  3. 节点C:AI模型评分(失败!因为模型对某个新网络用语无法判断,置信度极低)
  4. 节点D:最终处置(依赖C的评分)

如果没有 HITL Rerun,流程在C节点失败,整个任务就卡住了。管理员需要手动查看失败原因,然后用后台工具模拟C节点的输出,再手动触发D节点。这个过程繁琐且容易出错。

有了 HITL Rerun,引擎可以在C节点失败(或置信度低于阈值)时:

  1. 自动暂停流程,将C节点的输入、失败原因、以及上游A、B节点的结果,生成一个清晰的人工审核工单
  2. 工单分配给审核员。审核员查看内容,判断是否违规,并直接给出一个“人工评分”(替代C节点的输出)。
  3. 审核员提交后,引擎接收这个人工结果,将其作为C节点的成功输出,然后自动从C节点之后(即D节点)继续执行流程。

5.2 关键技术点:状态持久化与上下文恢复

实现 HITL Rerun 的核心挑战是状态管理。要能从某个节点重跑,引擎必须有能力:

  • 持久化执行状态:在流程执行到每一步时,将整个DAG的当前状态(哪些节点已完成及其输出、当前正在执行哪个节点、全局上下文数据)保存到数据库或分布式存储中。
  • 创建检查点:在可能需要进行人工干预的节点(如上述的C节点)之前,创建一个“检查点”。保存此刻之前所有节点的输出。
  • 注入人工结果:当人工处理完成后,系统需要能将人工提供的结果,准确地“注入”到对应节点的输出槽中,并标记该节点为“已完成(通过人工)”。
  • 从检查点恢复:引擎从存储中加载检查点的状态,用人工结果覆盖对应节点的输出,然后重新计算后续节点的依赖满足情况,并继续调度执行。

这要求ToolsNode的执行引擎不仅仅是内存中的调度器,还需要与一个状态持久化层紧密集成。每个ToolsNode的输入输出也最好是可序列化的,以便保存。

5.3 设计一个支持 HITL 的 ToolsNode

一个节点可以通过配置或继承来声明自己支持 HITL。

class HITLEnabledNode(ToolsNode): # 增加一个标志,表示该节点是否可能触发人工干预 requires_hitl_review: bool = False # 人工审核时的提示模板 hitl_prompt_template: str = "请审核以下内容:{input}, AI模型给出的置信度为{confidence}, 请给出最终判断。" async def execute(self, input_data: dict, context: ExecutionContext) -> dict: # ... 正常执行逻辑 ... result, confidence = await self._call_ai_model(input_data) if confidence < self.hitl_threshold: # 触发HITL流程 # 1. 抛出特定异常,或通过context设置状态 # 2. 引擎捕获后,会暂停流程,保存状态,创建工单 raise HumanInterventionRequired( node_name=self.name, input_data=input_data, intermediate_result={"ai_output": result, "confidence": confidence}, prompt=self.hitl_prompt_template.format(input=input_data, confidence=confidence) ) return {"final_result": result, "confidence": confidence}

执行引擎需要捕获HumanInterventionRequired异常,并触发后续的工单创建、状态保存流程。

实操心得:HITL 节点的设计哲学不要把 HITL 当作“万能兜底”。HITL 节点应该设计在确定性规则处理不了、但人工可以轻松判断的边界地带。同时,要尽可能为人工审核员提供丰富的上下文(上游节点结果、AI的中间推理、失败原因),降低其决策成本。每一次人工处理的结果,都应该考虑能否作为反馈数据,用于优化AI模型或调整规则,从而减少未来对HITL的依赖,实现闭环优化。

6. 实战:构建一个简易的 ToolsNode 引擎原型

理解了三大核心能力,我们动手设计一个极简的、具备这三方面特点的ToolsNode引擎原型,以加深理解。我们将使用 Python 的asyncio来实现并发。

6.1 定义核心类

首先,定义我们的ToolsNode基类和ExecutionContext

import asyncio import time from abc import ABC, abstractmethod from typing import Dict, List, Any, Callable, Optional from dataclasses import dataclass, field @dataclass class ExecutionContext: """执行上下文,用于在节点和中间件间传递全局数据""" request_id: str user_id: Optional[str] = None shared_data: Dict[str, Any] = field(default_factory=dict) # 全局共享数据袋 class ToolsNode(ABC): """工具节点抽象基类""" name: str description: str = "" @abstractmethod async def execute(self, input_data: Dict[str, Any], context: ExecutionContext) -> Dict[str, Any]: pass def get_dependencies(self) -> List[str]: """返回所依赖的节点名称列表,默认无依赖""" return []

6.2 实现中间件洋葱模型

实现一个支持中间件的执行器包装。

class Middleware: """中间件基类""" async def __call__(self, node: ToolsNode, input_data: Dict, context: ExecutionContext, next_fn: Callable): # 默认实现:直接调用下一个 return await next_fn(node, input_data, context) class LoggingMiddleware(Middleware): async def __call__(self, node, input_data, context, next_fn): start = time.time() print(f"[{context.request_id}] 进入节点: {node.name}, 输入: {input_data}") try: result = await next_fn(node, input_data, context) cost = time.time() - start print(f"[{context.request_id}] 节点成功: {node.name}, 耗时: {cost:.2f}s, 输出: {result}") return result except Exception as e: cost = time.time() - start print(f"[{context.request_id}] 节点失败: {node.name}, 耗时: {cost:.2f}s, 错误: {e}") raise class Orchestrator: """简单的执行编排器""" def __init__(self): self.nodes: Dict[str, ToolsNode] = {} self.middlewares: List[Middleware] = [] def register_node(self, node: ToolsNode): self.nodes[node.name] = node def add_middleware(self, middleware: Middleware): self.middlewares.append(middleware) def _wrap_with_middleware(self, node: ToolsNode) -> Callable: """将节点的execute方法用中间件层层包裹,形成洋葱结构""" async def final_executor(inp, ctx): return await node.execute(inp, ctx) # 从内到外包裹中间件 wrapped = final_executor for middleware in reversed(self.middlewares): # 注意顺序:先添加的中间件在外层 wrapped = (lambda m, n: lambda inp, ctx: m.__call__(n, inp, ctx, n))(middleware, node) # 简化写法,实际需闭包处理每个middleware # 为清晰起见,这里用一个简化版的包装逻辑: async def _execute(inp, ctx): # 构建中间件调用链 call_chain = final_executor for m in reversed(self.middlewares): call_chain = (lambda m, next_call: lambda i, c: m(node, i, c, next_call))(m, call_chain) return await call_chain(inp, ctx) return _execute

这段代码展示了洋葱模型的核心:通过高阶函数,将一个个中间件和最终的节点执行函数嵌套起来。实际项目中,可以使用starlettesentry-sdk等库中更成熟的中间件模式。

6.3 实现基于DAG的并行调度

现在,让编排器能够根据依赖关系调度节点。

async def execute_flow(self, start_nodes: List[str], initial_context: ExecutionContext, initial_data: Dict[str, Any]) -> Dict[str, Any]: """ 执行一个流程 :param start_nodes: 起始节点名列表 :param initial_context: 初始上下文 :param initial_data: 初始数据,key为节点名,value为输入 :return: 最终输出数据 """ from collections import deque, defaultdict # 1. 构建邻接表和入度表 adj = defaultdict(list) # 邻接表: node -> [下游节点] in_degree = defaultdict(int) # 节点入度 all_nodes = set(self.nodes.keys()) node_outputs = {} # 存储节点输出 node_inputs = {node_name: initial_data.get(node_name, {}) for node_name in all_nodes} # 存储节点输入(动态更新) # 初始化入度和邻接表 for node_name, node in self.nodes.items(): deps = node.get_dependencies() for dep in deps: if dep not in all_nodes: raise ValueError(f"节点 {node_name} 依赖了不存在的节点 {dep}") adj[dep].append(node_name) in_degree[node_name] += 1 # 2. 拓扑排序执行 (Kahn算法) queue = deque([n for n in start_nodes if in_degree[n] == 0]) executed_order = [] while queue: # 并行执行当前队列中所有可执行节点 current_batch = list(queue) queue.clear() # 清空队列,准备下一批 tasks = [] for node_name in current_batch: # 准备该节点的输入:依赖节点的输出 + 初始输入 input_data = node_inputs[node_name].copy() for dep in self.nodes[node_name].get_dependencies(): if dep in node_outputs: # 简单合并依赖输出,实际项目需更精细的映射 input_data.update(node_outputs[dep]) else: # 理论上不会发生,因为入度为0才执行 raise RuntimeError(f"依赖节点 {dep} 的输出未就绪") # 用中间件包裹后的执行函数 executor = self._wrap_with_middleware(self.nodes[node_name]) task = asyncio.create_task(executor(input_data, initial_context)) tasks.append((node_name, task)) # 等待这一批节点全部完成 results = await asyncio.gather(*[t for _, t in tasks], return_exceptions=True) # 处理结果,更新图状态 for (node_name, _), result in zip(tasks, results): if isinstance(result, Exception): # 错误处理:这里简单抛出,实际应更复杂 raise RuntimeError(f"节点 {node_name} 执行失败") from result node_outputs[node_name] = result executed_order.append(node_name) # 更新下游节点入度,并将新的可执行节点加入队列 for downstream in adj[node_name]: in_degree[downstream] -= 1 if in_degree[downstream] == 0: queue.append(downstream) if len(executed_order) != len(all_nodes): # 图中存在环,无法完全执行 remaining = [n for n in all_nodes if n not in executed_order] raise RuntimeError(f"检测到循环依赖,以下节点未执行: {remaining}") # 3. 返回最终输出(这里简单返回所有节点输出,实际可根据需要返回特定节点输出) return node_outputs

这个调度器实现了基本的拓扑排序和并行批量执行。它找出所有入度为0(没有未完成依赖)的节点,并行执行它们,然后更新依赖图,重复此过程。

6.4 模拟 HITL Rerun 的流程

我们在一个节点中模拟 HITL。为了简化,我们不实现完整的持久化工单系统,而是通过一个全局的“人工决策模拟器”来演示流程。

class HumanInterventionRequired(Exception): """触发人工干预的异常""" def __init__(self, node_name: str, input_data: Dict, context: ExecutionContext): self.node_name = node_name self.input_data = input_data self.context = context super().__init__(f"节点 {node_name} 需要人工干预") class AINodeWithLowConfidence(ToolsNode): """一个模拟的AI节点,置信度低时触发HITL""" def __init__(self, name): self.name = name self.description = "模拟AI节点,随机失败或低置信度" def get_dependencies(self): return ["start"] async def execute(self, input_data: Dict, context: ExecutionContext) -> Dict: import random # 模拟AI处理 await asyncio.sleep(0.5) confidence = random.random() # 0~1之间的随机数模拟置信度 if confidence < 0.3: # 置信度低于0.3,触发人工干预 print(f"\n[模拟HITL] 节点 {self.name} 置信度过低({confidence:.2f}), 触发人工审核。输入: {input_data}") # 在实际系统中,这里会抛异常,引擎捕获后创建工单、保存状态。 # 我们这里模拟人工决策过程: manual_decision = await self._simulate_human_review(input_data) return {"result": manual_decision, "confidence": 1.0, "source": "human"} elif confidence < 0.6: # 中等置信度,返回结果但标记 return {"result": f"AI结果(置信度{confidence:.2f})", "confidence": confidence, "source": "ai_low"} else: # 高置信度 return {"result": f"AI结果(置信度{confidence:.2f})", "confidence": confidence, "source": "ai_high"} async def _simulate_human_review(self, input_data): """模拟人工审核,等待2秒后返回一个确定结果""" print("[模拟HITL] 人工正在审核...") await asyncio.sleep(2) # 模拟人工总是批准 decision = f"人工审核通过: {input_data.get('query', '')}" print(f"[模拟HITL] 人工审核完成,决定: {decision}") return decision

6.5 运行一个完整示例

让我们把以上所有部分组合起来,运行一个包含并行执行、中间件和模拟HITL的小流程。

# 定义几个简单的节点 class StartNode(ToolsNode): def __init__(self): self.name = "start" self.description = "起始节点,处理用户输入" async def execute(self, input_data, context): await asyncio.sleep(0.2) query = input_data.get("query", "") return {"parsed_query": query, "location": "北京", "topic": "科技"} class WeatherNode(ToolsNode): def __init__(self): self.name = "weather" self.description = "查询天气" def get_dependencies(self): return ["start"] async def execute(self, input_data, context): await asyncio.sleep(1) # 模拟网络请求耗时 location = input_data.get("location", "") return {"weather": f"{location}天气晴朗,25度"} class NewsNode(ToolsNode): def __init__(self): self.name = "news" self.description = "检索新闻" def get_dependencies(self): return ["start"] async def execute(self, input_data, context): await asyncio.sleep(0.8) # 模拟另一个耗时请求 topic = input_data.get("topic", "") return {"news": f"关于{topic}的最新动态"} class ReportNode(ToolsNode): def __init__(self): self.name = "report" self.description = "生成最终报告" def get_dependencies(self): return ["weather", "news", "ai_node"] # 依赖三个节点 async def execute(self, input_data, context): # 合并所有依赖节点的输出 summary = f"天气:{input_data.get('weather')}。新闻:{input_data.get('news')}。AI分析:{input_data.get('result', 'N/A')}。" return {"final_report": summary} async def main(): # 1. 初始化编排器 orchestrator = Orchestrator() orchestrator.add_middleware(LoggingMiddleware()) # 2. 注册节点 orchestrator.register_node(StartNode()) orchestrator.register_node(WeatherNode()) orchestrator.register_node(NewsNode()) orchestrator.register_node(AINodeWithLowConfidence("ai_node")) orchestrator.register_node(ReportNode()) # 3. 准备执行 context = ExecutionContext(request_id="test_001", user_id="user_123") initial_data = {"start": {"query": "今天北京天气和科技新闻怎么样?"}} print("开始执行工具链...") try: # 4. 执行流程,从'start'节点开始 final_outputs = await orchestrator.execute_flow(start_nodes=["start"], initial_context=context, initial_data=initial_data) print("\n=== 执行完成 ===") for node_name, output in final_outputs.items(): print(f"{node_name}: {output}") print(f"\n最终报告: {final_outputs.get('report', {}).get('final_report', '无')}") except Exception as e: print(f"\n流程执行出错: {e}") if __name__ == "__main__": asyncio.run(main())

运行这个示例,你会看到:

  1. 日志中间件生效,每个节点的开始、结束、耗时都被打印。
  2. 并行执行weathernews节点会并行执行(因为它们都只依赖start),总耗时接近两者中较慢的那个(约1秒),而不是串行的1.8秒。
  3. 模拟HITLai_node有30%的概率因“置信度低”触发模拟的人工审核。你会看到相应的提示信息,并且流程会“等待”2秒模拟人工处理时间,然后继续。report节点会等待weather,news,ai_node全部完成后再执行。
  4. 依赖管理report节点正确等待了所有前置节点完成。

通过这个原型,我们亲手验证了ToolsNode三大核心能力如何在一个简易系统中协同工作。在实际的大型框架中(如 LangChain、AutoGPT 的底层设计,或企业内部的流程引擎),这些概念被实现得更加健壮、功能丰富,并集成了持久化、监控、版本管理等生产级特性。但万变不离其宗,其核心思想——通过声明式依赖实现并行、通过中间件实现管控、通过状态管理支持HITL——正是构建复杂、可靠、可观测的自动化工具链的基石。理解这些模式,能帮助我们在设计和选型时做出更明智的决策。

← 返回列表