十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

从零构建多Agent系统:Python实现自动化任务流水线

从零构建多Agent系统:Python实现自动化任务流水线 1. 先搞清楚 ZCODE 多 Agent 到底解决了什么问题如果你正在处理需要多个步骤、多种工具协作的复杂任务比如自动化的数据分析、跨平台的文件处理或者需要调用不同 API 来完成一个完整流程那么 ZCODE 多 Agent 这个方向就值得你停下来看看。它不是一个具体的工具而是一种架构思路核心是让多个具备特定能力的“智能体”协同工作把一个复杂的大任务拆解、分发、执行最后汇总结果。很多人一听到“多 Agent”就觉得是前沿学术概念离落地很远。但实际在工程化场景里它的价值非常直接把单点、手动的串行流程变成自动化、可并发的流水线。比如一个任务需要先爬取数据然后清洗接着调用模型分析最后生成报告。传统做法要么写一个巨长的脚本要么手动分步操作。而多 Agent 的思路是让“爬虫 Agent”、“清洗 Agent”、“分析 Agent”、“报告 Agent”各司其职通过一个调度中心来串联和监控。所以ZCODE 多 Agent 追求的更专业、更高效就体现在这里专业指的是每个 Agent 可以专注于自己最擅长的领域用最合适的工具和模型高效指的是通过并行、流水线、错误隔离和重试机制提升整体任务的成功率和吞吐量。这篇文章不会空谈理论我会结合常见的自动化开发场景拆解如何从零开始搭建一个可用的多 Agent 执行框架重点关注环境准备、Agent 定义、任务调度、错误处理这些真正影响落地的问题。2. 环境与核心组件别急着写代码先搭好架子在动手写第一个 Agent 之前得先把舞台搭好。多 Agent 系统不是魔法它建立在一些基础的开发环境和组件之上。我建议先从最小化的可行架构开始避免一开始就陷入复杂的框架选型。2.1 基础运行环境你的开发机就是最初的试验场。不需要多强的 GPU但需要稳定的网络和清晰的项目结构。编程语言与版本Python 是目前大多数 AI Agent 框架和工具链的首选。建议使用 Python 3.8-3.11 这些长期支持版本。用pyenv或conda创建独立的虚拟环境是必须的避免依赖冲突。# 创建并激活虚拟环境 python -m venv zcode_agent_env source zcode_agent_env/bin/activate # Linux/macOS # 或 zcode_agent_env\Scripts\activate # Windows核心依赖除了 Python 标准库有几个包几乎是必装的requests/aiohttp: 用于 Agent 与外部 API如模型服务、数据库通信。pydantic: 用于严格定义 Agent 之间传递的消息格式和数据模型这是保证协作不乱套的关键。redis/celery(可选): 如果你计划实现任务队列和异步执行这些是成熟的后台任务解决方案。但对于初期验证可以先用内存队列简化。项目目录结构一个清晰的目录结构能省去后期大量重构的麻烦。建议如下zcode_multi_agent/ ├── agents/ # 存放各个Agent的实现 │ ├── __init__.py │ ├── base_agent.py # Agent基类 │ ├── crawler_agent.py │ ├── analyzer_agent.py │ └── ... ├── core/ # 核心调度与通信逻辑 │ ├── __init__.py │ ├── orchestrator.py # 调度器/协调器 │ ├── message_bus.py # 消息总线初期可用内存实现 │ └── task.py # 任务定义 ├── tasks/ # 具体任务流程定义 ├── config/ # 配置文件 ├── logs/ # 日志目录 └── main.py # 入口文件2.2 理解多 Agent 系统的核心组件在写代码前脑子里要对这几个角色有清晰的认识Agent智能体系统的“工人”。每个 Agent 应该职责单一。例如工具调用型 Agent专门执行某个具体操作如调用搜索引擎 API、读写数据库、执行 Shell 命令。推理决策型 Agent基于 LLM大语言模型负责理解任务、做出判断、规划步骤。比如一个“任务规划 Agent”。专业处理型 Agent封装了特定领域能力如“PDF 解析 Agent”、“图像识别 Agent”。Orchestrator协调器/调度器系统的“大脑”或“项目经理”。它接收总任务将其分解成子任务分配给合适的 Agent并监控执行流程、处理异常。这是实现“高效”的关键。Message Bus消息总线Agent 之间的“通信管道”。Agent 不直接互相调用而是通过发布/订阅消息来协作。这解耦了各个组件使得系统更灵活、易于扩展。初期可以用一个简单的内存中的消息队列如queue.Queue或字典来模拟。Task Context任务与上下文需要定义清晰的数据结构来表示一个任务以及任务执行过程中的共享上下文Context。上下文里包含了初始输入、中间结果、最终输出等。3. 从零构建实现一个可运行的多 Agent 流水线现在我们用一个具体的例子串起来“获取技术文章摘要”。这个任务可以分解为1) 爬取指定文章正文2) 提取并总结核心内容3) 将摘要保存为 Markdown 文件。3.1 第一步定义基础数据模型在core/task.py和core/message.py中我们先定义 Agent 间通信的“语言”。# core/task.py from pydantic import BaseModel from typing import Any, Dict, Optional from enum import Enum class TaskStatus(str, Enum): PENDING pending RUNNING running SUCCESS success FAILED failed class Task(BaseModel): task_id: str name: str input_data: Dict[str, Any] # 任务输入如 {url: https://...} output_data: Optional[Dict[str, Any]] None # 任务输出 status: TaskStatus TaskStatus.PENDING assigned_agent: Optional[str] None # 被分配给哪个Agent# core/message.py from pydantic import BaseModel from typing import Any, Dict from .task import Task class AgentMessage(BaseModel): msg_id: str msg_type: str # 如 “task_assigned”, “task_result”, “error” sender: str # 发送方Agent ID receiver: str # 接收方Agent ID (或 “orchestrator”) payload: Dict[str, Any] # 消息内容如包含一个Task对象 timestamp: float3.2 第二步实现一个简单的 Agent 基类在agents/base_agent.py中创建一个所有 Agent 都继承的基类。# agents/base_agent.py import abc import logging from typing import Any, Dict from core.message import AgentMessage from core.task import Task, TaskStatus class BaseAgent(abc.ABC): def __init__(self, agent_id: str): self.agent_id agent_id self.logger logging.getLogger(self.agent_id) abc.abstractmethod async def execute(self, task: Task) - Task: 每个Agent必须实现的核心执行逻辑 pass def _mark_task_success(self, task: Task, result: Dict[str, Any]) - Task: task.status TaskStatus.SUCCESS task.output_data result self.logger.info(fTask {task.task_id} completed successfully.) return task def _mark_task_failed(self, task: Task, error: str) - Task: task.status TaskStatus.FAILED task.output_data {error: error} self.logger.error(fTask {task.task_id} failed: {error}) return task3.3 第三步实现具体的工作 Agent我们实现两个简单的 Agent一个爬虫一个总结器。这里为了简化爬虫用requests模拟总结用简单的文本截取模拟实际应接入 LLM。# agents/crawler_agent.py import requests from agents.base_agent import BaseAgent from core.task import Task import logging class CrawlerAgent(BaseAgent): def __init__(self, agent_id: str crawler_01): super().__init__(agent_id) async def execute(self, task: Task) - Task: self.logger.info(fStarting crawl task: {task.task_id}) url task.input_data.get(url) if not url: return self._mark_task_failed(task, No URL provided in input data.) try: # 实际项目中应添加headers、超时、重试等逻辑 response requests.get(url, timeout10) response.raise_for_status() # 假设我们只获取文本内容 article_content response.text[:5000] # 限制长度做演示 result {raw_content: article_content, source_url: url} return self._mark_task_success(task, result) except requests.RequestException as e: return self._mark_task_failed(task, fFailed to fetch URL: {e})# agents/summary_agent.py from agents.base_agent import BaseAgent from core.task import Task class SummaryAgent(BaseAgent): def __init__(self, agent_id: str summarizer_01): super().__init__(agent_id) async def execute(self, task: Task) - Task: self.logger.info(fStarting summary task: {task.task_id}) raw_content task.input_data.get(raw_content) if not raw_content: return self._mark_task_failed(task, No content to summarize.) # 模拟一个极其简单的“总结”取前200字符作为摘要 # 真实场景这里应调用LLM API如OpenAI, DeepSeek, 或本地模型 simulated_summary raw_content[:200] ... [Summarized] result {summary: simulated_summary, original_length: len(raw_content)} return self._mark_task_success(task, result)3.4 第四步实现一个简单的协调器协调器负责接收主任务创建子任务并驱动 Agent 工作。这里实现一个最简单的同步版本。# core/orchestrator.py import asyncio import logging from typing import Dict, List from core.task import Task, TaskStatus from core.message import AgentMessage from agents.crawler_agent import CrawlerAgent from agents.summary_agent import SummaryAgent class SimpleOrchestrator: def __init__(self): self.logger logging.getLogger(orchestrator) self.agents { crawler: CrawlerAgent(), summarizer: SummaryAgent(), } self.task_registry: Dict[str, Task] {} async def process_article_summary(self, url: str) - Dict: 处理文章摘要的主流程 main_task_id fmain_{hash(url)} main_task Task(task_idmain_task_id, namearticle_summary, input_data{url: url}) self.task_registry[main_task_id] main_task self.logger.info(fStarted main task: {main_task_id} for URL: {url}) # 步骤1爬取 crawl_task Task(task_idf{main_task_id}_step1, namecrawl_article, input_data{url: url}) self.task_registry[crawl_task.task_id] crawl_task crawl_task await self.agents[crawler].execute(crawl_task) if crawl_task.status ! TaskStatus.SUCCESS: main_task.status TaskStatus.FAILED main_task.output_data {error: fCrawl failed: {crawl_task.output_data.get(error)}} return main_task.dict() # 步骤2总结 summary_task Task( task_idf{main_task_id}_step2, namesummarize_content, input_data{raw_content: crawl_task.output_data[raw_content]} ) self.task_registry[summary_task.task_id] summary_task summary_task await self.agents[summarizer].execute(summary_task) if summary_task.status ! TaskStatus.SUCCESS: main_task.status TaskStatus.FAILED main_task.output_data {error: fSummary failed: {summary_task.output_data.get(error)}} return main_task.dict() # 汇总结果 main_task.status TaskStatus.SUCCESS main_task.output_data { original_url: url, summary: summary_task.output_data[summary], crawl_status: success, summary_status: success } self.logger.info(fMain task {main_task_id} completed successfully.) return main_task.dict()3.5 第五步创建入口并运行测试在项目根目录创建main.py。# main.py import asyncio import logging from core.orchestrator import SimpleOrchestrator # 配置日志方便查看执行过程 logging.basicConfig(levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s) async def main(): orchestrator SimpleOrchestrator() # 用一个示例URL测试 test_url https://httpbin.org/html # 一个返回示例HTML的测试网站 result await orchestrator.process_article_summary(test_url) print(\n 任务执行结果 ) print(f最终状态: {result[status]}) if result[status] success: print(f文章摘要: {result[output_data][summary]}) else: print(f错误信息: {result[output_data]}) if __name__ __main__: asyncio.run(main())运行这个脚本 (python main.py)你会在控制台看到两个 Agent 依次被调用并最终输出摘要结果。虽然简单但这已经是一个完整的多 Agent 工作流水线。4. 从“能跑”到“高效专业”关键优化与避坑点上面的例子只是一个起点。要让多 Agent 系统真正变得“专业高效”必须解决以下几个工程问题。这也是新手最容易踩坑的地方。4.1 通信机制从同步到异步引入消息队列我们上面的协调器是同步等待每个 Agent 完成的这没有发挥多 Agent 的并发优势。生产环境中Agent 应该是独立运行的通过消息队列通信。为什么需要消息队列实现解耦和异步。协调器发布一个爬虫任务后不用等待可以继续处理其他逻辑。爬虫 Agent 完成任务后将结果发回消息队列再由协调器或下一个 Agent 消费。如何做可以使用RabbitMQ、Redis Streams或Kafka。对于 Pythonceleryredis是经典组合。协调器将任务推入队列每个 Agent 作为一个独立的 Celery Worker 从队列中拉取任务执行。避坑消息格式必须统一且版本兼容。使用像Pydantic这样的库严格序列化/反序列化消息体避免因字段不一致导致解析失败。4.2 任务调度与依赖管理复杂任务中子任务间可能有依赖关系。比如“总结 Agent” 必须等 “爬虫 Agent” 完成后才能开始。实现 DAG有向无环图可以使用Airflow、Prefect这类调度框架来定义任务流。也可以自己实现一个简单的状态机在协调器中维护任务状态图。关键点每个任务完成后要能触发其下游依赖任务的检查。如果下游任务的所有前置任务都完成了就将其放入就绪队列等待调度。避坑小心循环依赖。在设计任务流时必须确保它是“无环”的。可以用拓扑排序算法在提交前做检查。4.3 错误处理与重试机制这是保证系统鲁棒性的核心。一个 Agent 失败不能导致整个流程崩溃。分级错误处理瞬时错误如网络超时、API 限流应自动重试并采用指数退避策略。业务逻辑错误如输入格式不对、资源不存在应标记任务失败记录明确错误原因并可能触发告警。系统级错误如依赖服务宕机、内存溢出应停止调度新任务并向上游报告系统不可用。实现重试可以在 Agent 的execute方法内部包装重试逻辑也可以由协调器在收到失败消息后重新派发任务需注意幂等性。# 在Agent基类或工具函数中实现带退避的重试 import time async def execute_with_retry(agent, task, max_retries3): for attempt in range(max_retries): try: return await agent.execute(task) except TemporaryError as e: # 自定义的瞬时异常 if attempt max_retries - 1: raise wait_time (2 ** attempt) (random.random() * 0.1) # 指数退避加抖动 time.sleep(wait_time) continue死信队列对于重试多次仍失败的任务应将其移入死信队列供人工或更高级的修复流程处理。4.4 资源管理与限流当你有大量并发任务时必须考虑资源限制。限制并发数为每个 Agent 类型设置最大并发 Worker 数。例如同时运行的“爬虫 Agent”不能超过 5 个避免把目标网站爬崩或耗尽本地连接池。资源隔离如果 Agent 调用的是消耗 GPU 的模型更需要严格控制并发甚至使用专门的队列和机器。实现方式在协调器派发任务时检查该类型 Agent 的“正在运行任务数”。也可以利用消息队列如 RabbitMQ的 QoS 设置或 Celery 的worker_concurrency参数。4.5 可观测性日志、监控与追踪系统复杂后必须能看清里面发生了什么。结构化日志不要简单用print。使用logging模块输出 JSON 格式的日志方便后续用 ELKElasticsearch, Logstash, Kibana或 Loki 收集和查询。在日志中务必包含task_id,agent_id,step等关键上下文。分布式追踪为每个主任务生成一个唯一的trace_id并在这个任务流经的所有 Agent 和组件中传递这个 ID。这样可以在日志和监控系统中轻松串联起一个请求的完整生命周期。可以使用OpenTelemetry这样的标准。关键指标监控任务吞吐量TPS任务平均耗时、P95/P99 耗时各 Agent 的成功率、失败率队列长度系统资源使用率CPU、内存、网络 这些指标可以通过 Prometheus 暴露用 Grafana 展示。5. 进阶思考何时需要以及如何设计更复杂的 Agent并不是所有任务都需要多 Agent。一个简单的脚本能搞定的事情上多 Agent 就是过度设计。5.1 适用多 Agent 的场景特征流程长且步骤分明任务可以清晰地拆解为多个顺序或并行的阶段。需要多种专业能力不同步骤需要完全不同的工具、模型或外部服务。对容错和扩展性有要求希望某个环节失败不影响整体或希望轻松扩展某个环节的处理能力。需要动态规划任务路径不是完全固定的需要根据中间结果动态决定下一步。5.2 设计“更智能”的 Agent我们之前的 Agent 是“工具人”被动接收指令。更高级的 Agent 可以具备一定自主性。规划型 Agent输入一个模糊目标如“分析某公司竞品”由这个 Agent 调用 LLM自主拆解出具体步骤搜集信息、整理数据、对比分析、生成报告并生成一个可执行的任务图交给协调器。工具调用型 Agent给 Agent 配备一个“工具包”函数列表并赋予它使用 LLM 理解用户指令、自主选择并调用工具的能力。这就是类似 AutoGPT、ChatGPT Plugins 的思路。实现时需要让 LLM 输出结构化的工具调用请求如 JSON然后由 Agent 的执行器去解析并执行。多轮对话与记忆让 Agent 具备短期记忆对话历史和长期记忆向量数据库使其能在多轮交互中保持一致性和上下文理解。5.3 与现有架构集成作为微服务将每个 Agent 打包成独立的 HTTP 或 gRPC 服务通过 API 网关进行调用。这样可以用不同语言实现不同 Agent也便于独立部署和伸缩。集成到现有数据管道将多 Agent 系统作为现有 ETL提取、转换、加载管道或数据处理平台中的一个智能环节。从 Kafka 消费数据处理后再推回 Kafka。人机协同设计一些“人工审核 Agent”或“人工干预点”。当自动处理置信度低或遇到特定错误时将任务挂起并通知人工处理。构建一个专业高效的多 Agent 系统核心不在于追求 Agent 的“智能”程度而在于架构的清晰、通信的可靠、错误的可处理以及状态的可观测。先从一个小而具体的流水线跑起来然后逐步引入消息队列、错误重试、资源限流和监控。当你发现手动串联脚本已经变得难以维护和扩展时就是考虑 ZCODE 多 Agent 这类思路的最佳时机。记住合适的工具用在合适的复杂度上才是最高效的做法。
返回列表