大模型实战 07 额外篇] 从 ReAct 到 Workflow:基于云端 API 构建事件驱动的智能体
·
大模型实战 07 额外篇:从 ReAct 到 Workflow:基于云端 API 构建事件驱动的智能体
引言:从简单问答到智能工作流在之前的文章中,我们探讨了 ReAct(Reasoning + Acting)模式,它让大模型能够通过“思考-行动-观察”的循环来解决复杂问题。然而,在实际生产环境中,单纯的 ReAct 模式往往面临效率低、状态管理复杂、无法处理长时间运行任务等问题。今天,我们将深入探讨如何从 ReAct 升级到 事件驱动的 Workflow 架构,并基于云端 API 构建一个真正可落地、可扩展的智能体系统。## 为什么需要 Workflow?ReAct 的核心是“模型驱动循环”:模型思考一步,执行一步,再观察结果。这在简单场景下工作良好,但有几个痛点:1. 阻塞式执行:模型必须等待每一步的API返回结果,无法并行处理。2. 无状态管理:每次对话上下文都包含完整历史,导致Token浪费和大模型幻觉。3. 错误处理困难:如果某一步API调用超时或失败,整个循环可能卡死。Workflow 模式将任务拆解为多个独立的事件节点,通过消息队列和状态机来编排。每个节点可以异步触发,由事件总线驱动。这样,智能体不再是“思考-执行”的线性流程,而是一个可扩展的事件处理引擎。## 核心架构:事件驱动智能体我们将构建一个基于云端API的事件驱动智能体,核心组件包括:- 事件总线(Event Bus):使用Redis Pub/Sub或Kafka- 状态存储:Redis或PostgreSQL存储任务状态- 大模型节点:调用OpenAI或Claude API进行推理- 工具节点:调用外部API(如天气、数据库、文件系统)- 编排器(Orchestrator):定义工作流DAG下面是我们架构的简化图示(文字描述):用户输入 → 事件总线 → 编排器解析意图 → 拆分任务 → 触发多个事件节点 ↓ 大模型推理节点 → 产生新事件 ↓ 工具执行节点 → 返回结果 ↓ 聚合节点 → 回复用户## 实战代码:从 ReAct 到 Workflow 的演变### 示例1:基于 ReAct 的简单智能体(对比基准)我们先写一个基础的 ReAct 智能体,用于查询天气并计算温差。python# react_agent.py - 基础ReAct实现import jsonimport requestsfrom openai import OpenAIclass ReactAgent: def __init__(self, api_key): self.client = OpenAI(api_key=api_key) self.history = [] def think_and_act(self, query): """ReAct主循环:思考→行动→观察""" system_prompt = """你是一个智能助手,可以调用以下工具: - get_weather(city): 获取城市天气 - calculate_difference(a, b): 计算两个数的差 请用JSON格式输出你的思考和行动。""" self.history.append({"role": "user", "content": query}) for step in range(5): # 最多5步 # 1. 思考 response = self.client.chat.completions.create( model="gpt-4", messages=[ {"role": "system", "content": system_prompt}, *self.history ] ) thought = response.choices[0].message.content self.history.append({"role": "assistant", "content": thought}) # 2. 解析行动 try: action = json.loads(thought) if action.get("action") == "get_weather": result = self._get_weather(action["city"]) elif action.get("action") == "calculate_difference": result = self._calc_diff(action["a"], action["b"]) else: break # 没有行动,结束 except: break # JSON解析失败,结束 # 3. 观察结果 observation = f"工具返回: {result}" self.history.append({"role": "user", "content": observation}) # 如果结果包含答案,提前结束 if "最终答案" in thought: return thought return self.history[-1]["content"] def _get_weather(self, city): # 模拟天气API调用 weather_data = {"北京": {"temp": 25, "humidity": 60}, "上海": {"temp": 30, "humidity": 70}} return weather_data.get(city, {"temp": 20, "humidity": 50}) def _calc_diff(self, a, b): return abs(a - b)# 使用示例if __name__ == "__main__": agent = ReactAgent(api_key="your_openai_key") result = agent.think_and_act("北京和上海的温差是多少?") print(result)这段代码的问题:所有步骤串行执行,如果get_weather调用耗时3秒,整个对话就被阻塞。而且历史记录越来越长,Token消耗大。### 示例2:基于 Workflow 的事件驱动智能体现在,我们用事件驱动的方式重构。这里我们使用 Redis 作为事件总线,结合简单的状态机。python# workflow_agent.py - 事件驱动智能体import redisimport jsonimport asyncioimport aiohttpfrom datetime import datetimeclass EventDrivenAgent: def __init__(self, redis_host="localhost", redis_port=6379): self.redis = redis.Redis(host=redis_host, port=redis_port, decode_responses=True) self.pubsub = self.redis.pubsub() async def start_workflow(self, user_query: str): """启动新的工作流""" workflow_id = f"wf_{datetime.now().timestamp()}" # 1. 初始化工作流状态 workflow_state = { "id": workflow_id, "query": user_query, "status": "init", "steps": [], "results": {} } self.redis.set(f"workflow:{workflow_id}", json.dumps(workflow_state)) # 2. 发布意图解析事件 event = { "type": "intent_parse", "workflow_id": workflow_id, "data": {"query": user_query} } self.redis.publish("agent_events", json.dumps(event)) # 3. 监听结果事件(异步) asyncio.create_task(self._listen_for_result(workflow_id)) return workflow_id async def _listen_for_result(self, workflow_id): """监听工作流完成事件""" self.pubsub.subscribe(f"result:{workflow_id}") for message in self.pubsub.listen(): if message["type"] == "message": data = json.loads(message["data"]) if data["status"] == "completed": print(f"工作流 {workflow_id} 完成: {data['final_result']}") return data['final_result'] async def intent_parse_handler(self, event): """处理意图解析(事件处理器)""" workflow_id = event["workflow_id"] query = event["data"]["query"] # 调用大模型解析意图(这里简化) llm_response = await self._call_llm(f"解析以下查询的意图和工具调用: {query}") # 更新状态 state = json.loads(self.redis.get(f"workflow:{workflow_id}")) state["steps"].append({"step": "intent_parse", "result": llm_response}) self.redis.set(f"workflow:{workflow_id}", json.dumps(state)) # 根据意图触发不同事件 if "天气" in query: await self._trigger_weather_events(workflow_id, query) else: # 默认触发通用处理 event = { "type": "general_task", "workflow_id": workflow_id, "data": {"query": query} } self.redis.publish("agent_events", json.dumps(event)) async def _trigger_weather_events(self, workflow_id, query): """触发多个并行的天气查询事件""" # 并行查询北京和上海的天气 cities = ["北京", "上海"] for city in cities: event = { "type": "weather_query", "workflow_id": workflow_id, "data": {"city": city} } self.redis.publish("agent_events", json.dumps(event)) async def weather_query_handler(self, event): """处理天气查询事件""" workflow_id = event["workflow_id"] city = event["data"]["city"] # 模拟异步API调用 async with aiohttp.ClientSession() as session: # 实际中替换为真实天气API async with session.get(f"http://api.weather.com/{city}") as resp: weather_data = await resp.json() # 发布天气结果事件 result_event = { "type": "weather_result", "workflow_id": workflow_id, "data": {"city": city, "weather": weather_data} } self.redis.publish("agent_events", json.dumps(result_event)) # 检查是否两个城市都查询完毕 state = json.loads(self.redis.get(f"workflow:{workflow_id}")) state["results"][city] = weather_data if len(state["results"]) == 2: # 两个城市都返回了 # 触发计算温差事件 event = { "type": "calculate_diff", "workflow_id": workflow_id, "data": {"cities": list(state["results"].keys())} } self.redis.publish("agent_events", json.dumps(event)) self.redis.set(f"workflow:{workflow_id}", json.dumps(state)) async def calculate_diff_handler(self, event): """处理温差计算事件""" workflow_id = event["workflow_id"] state = json.loads(self.redis.get(f"workflow:{workflow_id}")) temps = [state["results"][city]["temp"] for city in event["data"]["cities"]] diff = abs(temps[0] - temps[1]) # 发布最终结果 final_result = f"温差为: {diff}度" self.redis.publish(f"result:{workflow_id}", json.dumps({ "status": "completed", "final_result": final_result })) async def _call_llm(self, prompt): """调用大模型API(异步)""" # 实际中调用OpenAI/Claude API await asyncio.sleep(0.5) # 模拟延迟 return {"intent": "weather_query", "parameters": {"city": "北京,上海"}} async def run_event_loop(self): """主事件循环""" self.pubsub.subscribe("agent_events") # 注册事件处理器 handlers = { "intent_parse": self.intent_parse_handler, "weather_query": self.weather_query_handler, "calculate_diff": self.calculate_diff_handler, } print("事件驱动智能体启动...") for message in self.pubsub.listen(): if message["type"] == "message": event = json.loads(message["data"]) event_type = event["type"] if event_type in handlers: # 异步处理事件,不阻塞主循环 asyncio.create_task(handlers[event_type](event))# 使用示例if __name__ == "__main__": agent = EventDrivenAgent() async def main(): # 启动事件循环 loop = asyncio.create_task(agent.run_event_loop()) # 模拟用户输入 workflow_id = await agent.start_workflow("北京和上海的温差是多少?") print(f"工作流已启动: {workflow_id}") # 等待工作流完成 await asyncio.sleep(5) asyncio.run(main())这段代码的亮点:1. 异步非阻塞:使用asyncio和aiohttp,天气查询可以并行执行。2. 事件驱动:通过Redis Pub/Sub解耦各组件,可以轻松添加新的事件处理器。3. 状态持久化:工作流状态存储在Redis中,支持故障恢复和横向扩展。4. 可组合性:每个事件处理器是独立单元,可以单独测试和部署。## 云端部署与扩展在实际生产环境中,我们还需要考虑:1. 使用消息队列:用Kafka或RabbitMQ替代Redis Pub/Sub,支持消息持久化和重试。2. 容器化部署:用Docker打包各个事件处理器,通过Kubernetes编排。3. 监控与日志:集成OpenTelemetry追踪事件流。4. 错误处理:实现死信队列和重试机制。下面是一个简单的Docker Compose配置示例:yaml# docker-compose.ymlversion: '3.8'services: redis: image: redis:7-alpine ports: - "6379:6379" agent-orchestrator: build: . depends_on: - redis environment: - REDIS_HOST=redis command: python orchestrator.py weather-worker: build: . depends_on: - redis environment: - REDIS_HOST=redis command: python weather_worker.py llm-worker: build: . depends_on: - redis environment: - OPENAI_API_KEY=${OPENAI_API_KEY} - REDIS_HOST=redis command: python llm_worker.py## 总结从 ReAct 到 Workflow 的演进,本质上是智能体架构从单线程思考器向分布式事件处理系统的升级。通过本文的实战代码,我们看到了:- ReAct 模式适合快速原型,但难以应对复杂、长时间运行的任务。- 事件驱动 Workflow 通过解耦和异步处理,实现了更高的吞吐量和可扩展性。- 基于云端 API(Redis + aiohttp + OpenAI)的架构,可以轻松集成到现有微服务体系中。关键收获:智能体不应该是一个“会思考的循环”,而应该是一个事件编排系统——大模型只是其中的一个处理器节点,与其他工具节点地位平等。这种架构让智能体真正具备了“工业化”的能力:可扩展、可容错、可观测。下一步,你可以尝试将本文的示例扩展到更复杂的场景,比如结合RAG(检索增强生成)做文档问答,或者集成Slack/飞书等消息平台,构建一个真正的生产级智能助手。
更多推荐

所有评论(0)