ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

微信调度AI Agent实战:Python构建聊天式任务调度系统

微信调度AI Agent实战:Python构建聊天式任务调度系统 直接用微信发消息就能让 AI 智能体去执行任务这听起来像是科幻片里的场景但现在已经有不少开发者在借助社交软件作为 Agent 的调度入口。这篇文章会从实际需求出发讲清楚为什么要用社交软件来调度 AI Agent、如何设计一套可运行的调度架构、以及完整的代码实现思路。整套方案基于 Python 实现不需要自建重型的任务中间件适合中小型个人项目或团队内部工具使用。1. 为什么需要用微信这类社交软件来调度 Agent先把一个关键问题放在前面市面上已经有那么多 Agent 开发框架也有一堆任务调度工具为什么还需要用微信来调度 AI 智能体看一个真实场景。假设你给团队搭了一个内部 AI Agent它能够根据项目需求自动搜索资料、梳理竞品信息、生成周报。你现在的使用方式是什么打开电脑登录某个管理后台进入 Agent 页面点击执行按钮。听起来不算繁琐但一旦任务多起来就会发现所有操作都被绑定在固定的终端上。另外一个更常见的情况是你正在通勤路上突然想到一个问题需要 Agent 去处理。因为 Agent 部署在服务器上你没法直接操作只能等回到电脑前。这时如果 Agent 能够挂在微信公众号、企业微信、钉钉或者飞书上你可以直接像聊天一样发一条指令帮我整理本周的竞品动态重点关注移动端产品的更新。消息发送到 AgentAgent 开始执行任务完成后把结果推回聊天窗口。整个过程不再依赖固定设备也不需要额外安装一套复杂的客户端。这就是用社交软件调度 Agent 的核心价值把Agent 的使用门槛降低到“发消息”这个层面。再从工程视角看社交软件本质上是一个天然的异步消息系统。消息的收发、会话状态、用户身份认证等基础能力都已经很成熟不需要自己从零实现。我们真正要做的事情只有三件接收社交软件推送过来的用户消息。解析消息内容识别用户意图确定要调用哪个 Agent。把 Agent 的执行结果回传给用户。这三件事对应到技术架构上分别就是消息接入层、任务调度层和结果回传层。后续的所有代码和配置都会围绕这个三层结构展开。2. Agent 调度的基础概念与整体架构在进入代码之前先把几个容易混淆的概念理清楚。很多读者看到“Agent 调度”第一反应是 K8s 里的 Pod 调度或者是分布式系统里的任务分发。实际上用社交软件调度 Agent和这些底层调度完全是两回事。2.1 Agent 是什么Agent 是一个能根据目标自动拆分任务、调用工具、决定执行步骤的智能程序。你可以把它理解成一个“有手有脚”的大模型大模型负责思考和生成Agent 则把它和环境连接起来让它能够查数据库、调 API、操作浏览器、读写文件。在本文的架构里Agent 不一定是完整的多步骤智能体也可以是一个封装好的函数、一个自动化脚本、或者一个调用大模型 API 的 Worker。对调度系统来说每个 Agent 只需要暴露一个统一的执行接口即可。2.2 什么是调度调度在这里的含义是“根据用户消息匹配并触发对应的 Agent 执行”。它不涉及 Kubernetes、容器或服务编排更像是一个轻量级的任务路由层。调度逻辑需要处理的关键点包括消息多义性同一句话可能对应多个 Agent需要识别意图优先级。任务状态管理Agent 执行需要时间消息收到后要记录任务状态不能让用户以为消息“石沉大海”。并发控制多个用户同时发消息时不能让所有请求直接涌入同一个 Agent造成资源竞争。2.3 整体架构分层整个系统可以拆成四层层级职责技术方案消息接入层接收微信/企微/钉钉/飞书的回调消息校验签名ItChat、微信公众平台开发接口、企业微信应用消息回调意图识别层解析用户文本匹配 Agent 名称和参数正则匹配 LLM 函数调用轻量场景用关键词规则即可任务调度层分发任务、维护状态、执行结果采集进程内队列 SQLite 持久化记录结果回传层将结果格式化成消息通过接入层回复用户文本回复 富文本Markdown/卡片这四层结构对新手来说足够清晰。如果你的团队已经引入 MQ 或消息中间件也可以把任务调度层替换成独立消息队列否则进程内队列和数据库状态表反而是更快捷的落地方式。3. 环境准备与前置条件为了让后面的代码可以实际跑通需要提前准备一些环境。这里以 Windows / macOS / Linux 均可运行为准版本说明以“通用版本”为指导不固定绑定某个具体版本号。3.1 运行环境与依赖基础环境建议Python 3.9 及以上版本。能够访问公网如果使用云服务器部署需要一个公网回调地址。需要一个可用的社交软件开发者账号或消息接口。需要安装的核心依赖库pip install itchat-uos requests pydantic flask说明一下依赖用途库名用途itchat-uos个人微信消息收发接入仅用于个人玩具项目或内网测试flask启动一个轻量 Web 服务接收第三方回调消息requests调用大模型 API、发送 HTTP 请求pydantic配置解析与消息结构校验使用 itchat-uos 需要特别提醒该库基于个人微信网页版协议存在被官方限制的风险不适合生产环境。如果是企业内部使用建议优先考虑企业微信的“自建应用消息回调”下面代码会同时给出企业微信回调接口的接收逻辑。3.2 大模型 API KeyAgent 的核心能力来自 LLM 调用。这里需要申请一个大模型服务的 API Key。目前常见的方案包括 OpenAI API、国内大模型平台的兼容接口等。在你的环境变量中配置export LLM_API_KEYyour-api-key export LLM_API_BASEhttps://api.example.com export LLM_MODELyour-model-name如果只是先做功能验证也可以不依赖大模型用一个最简单的关键词匹配逻辑充当代 Agent。后面示例代码里会体现这种可替换性。4. 消息接收层打通微信 / 企业微信消息入口消息接收是整个系统的基础。日常使用中“微信调度 Agent”有两种主流实现方式基于个人微信协议实现如 itchat 系列库。基于公众号 / 企业微信官方回调实现。4.1 个人微信接收消息演示用个人微信协议方案的问题在于接口限制和稳定性。它适合个人开发者在测试环境验证流程不适合正式对内服务。示例代码如下# 文件路径agent_scheduler/message_handlers/wechat_personal.py import itchat from itchat.content import TEXT from core.dispatcher import dispatch_message itchat.msg_register(TEXT) def text_reply(msg): # 把收到的文本消息交给调度器处理 user_id msg.user.userName content msg.text.strip() reply dispatch_message(user_id, content) return reply def start_wechat_bot(): # 启动个人微信消息监听 itchat.auto_login(hotReloadTrue) itchat.run()需要注意dispatch_message是后续调度层的统一入口。设计这个入口时要保证消息来源无关即不管是微信、钉钉还是飞书来的消息最终都走同一个函数。这样方便后续扩展。4.2 企业微信 / 公众号官方回调推荐方案更可靠的方案是通过企业微信自建应用接收消息。核心逻辑是启动一个 Flask 服务处理企业微信服务器发来的消息回调。# 文件路径agent_scheduler/message_handlers/wecom.py import hashlib import xml.etree.ElementTree as ET from flask import Flask, request, make_response from core.dispatcher import dispatch_message app Flask(__name__) # 企业微信配置 CORP_ID your-corpid AGENT_ID your-agentid SECRET your-secret TOKEN your-token ENCODING_AES_KEY your-encoding-aes-key app.route(/wecom/callback, methods[GET, POST]) def wecom_callback(): # 签名校验防止伪造请求 signature request.args.get(msg_signature) timestamp request.args.get(timestamp) nonce request.args.get(nonce) echostr request.args.get(echostr) # 简化处理实际项目需要按企业微信文档校验签名并解密 if request.method GET: # URL 验证 return make_response(echostr, 200) # POST 消息体是 XML 格式 xml_data request.data.decode(utf-8) root ET.fromstring(xml_data) msg_type root.find(MsgType).text content root.find(Content).text if msg_type text else 非文本消息暂不支持 user_id root.find(FromUserName).text # 交给统一的调度入口 reply dispatch_message(user_id, content) # 构造被动回复的 XML简化示例 resp_xml fxml ToUserName![CDATA[{user_id}]]/ToUserName FromUserName![CDATA[{AGENT_ID}]]/FromUserName CreateTime1234567890/CreateTime MsgType![CDATA[text]]/MsgType Content![CDATA[{reply}]]/Content /xml return make_response(resp_xml, 200) if __name__ __main__: app.run(host0.0.0.0, port8000, debugFalse)这段代码的核心意义在于演示回调接口的标准姿态接收请求、校验身份、解析内容、交给调度入口、返回结果。生产环境里这三个环节都需要认真加固。5. 统一消息分发实现意图识别与 Agent 匹配所有消息最终都会进入dispatch_message。这一步要解决两个问题用户想做什么。把任务交给哪个 Agent。5.1 基于注册表的 Agent 管理先设计一个 Agent 注册表。每个 Agent 都包含名称、描述、关键词、执行函数。注册表的好处是后续新增 Agent 只需追加一条配置不需要改动调度逻辑。# 文件路径agent_scheduler/core/agent_registry.py from dataclasses import dataclass from typing import Callable, Dict, List dataclass class AgentSpec: name: str description: str keywords: List[str] handler: Callable[[str, str], str] class AgentRegistry: def __init__(self): self._agents: Dict[str, AgentSpec] {} def register(self, spec: AgentSpec): self._agents[spec.name] spec print(f[AgentRegistry] 已注册 Agent: {spec.name}) def match(self, message: str) - AgentSpec | None: 根据消息内容匹配最合适的 Agent。 匹配策略优先精确匹配 Agent 名称其次按关键词匹配。 for spec in self._agents.values(): if message.startswith(f{spec.name}): return spec for keyword in spec.keywords: if keyword in message: return spec return None def list_agents(self) - List[str]: return list(self._agents.keys()) # 全局注册表实例 registry AgentRegistry()这里没有引入复杂的 NLP 模型原因是对于“Agent名 指令文本”的调度场景关键词匹配已经能覆盖 90% 的需求。如果要在真实项目中支持更复杂的语义可以在match方法里接入 LLM 的意图分类接口。5.2 调度入口实现消息分发入口的职责是先匹配 Agent再启动执行。这里的执行不是同步阻塞的而是把任务放入队列返回“已收到”的提示。# 文件路径agent_scheduler/core/dispatcher.py import uuid from datetime import datetime from core.agent_registry import registry from core.task_store import TaskStore from core.worker import TaskWorker task_store TaskStore() worker TaskWorker() def dispatch_message(user_id: str, content: str) - str: 统一消息调度入口。 :param user_id: 用户唯一标识 :param content: 用户发送的文本消息 :return: 需要回复给用户的文本 # 1. 生成任务 ID task_id uuid.uuid4().hex[:12] # 2. 匹配 Agent agent registry.match(content) if agent is None: return 无法识别意图。可用 Agent , .join(registry.list_agents()) # 3. 保存任务状态 task_store.create_task( task_idtask_id, user_iduser_id, agent_nameagent.name, raw_messagecontent, statusPENDING, created_atdatetime.now().isoformat(), ) # 4. 提交给后台 Worker 异步执行 worker.submit(task_idtask_id, agentagent, user_iduser_id, messagecontent) # 5. 立即返回提示 return f任务已受理任务号{task_id}Agent「{agent.name}」正在处理完成后会通知你。这个入口函数写得比较稳的一个原因是任务状态和实际执行分离。用户发送消息后能立刻得到反馈Agent 在后台运行完再回传结果。如果你做成同步等待模型推理完成再回复用户可能要等几十秒体验会差很多。6. 任务状态管理与异步执行 Worker调度系统没有任务状态管理就像没有日志的定时任务一样出了问题完全没法排查。这里用 SQLite 做状态存储结构简单又足够可靠。6.1 任务状态表# 文件路径agent_scheduler/core/task_store.py import sqlite3 import json from datetime import datetime class TaskStore: def __init__(self, db_pathagent_tasks.db): self._db_path db_path self._init_db() def _get_conn(self): conn sqlite3.connect(self._db_path) conn.row_factory sqlite3.Row return conn def _init_db(self): with self._get_conn() as conn: conn.execute( CREATE TABLE IF NOT EXISTS tasks ( task_id TEXT PRIMARY KEY, user_id TEXT NOT NULL, agent_name TEXT NOT NULL, raw_message TEXT, status TEXT NOT NULL, result TEXT, created_at TEXT, updated_at TEXT ) ) def create_task(self, task_id: str, user_id: str, agent_name: str, raw_message: str, status: str, created_at: str): with self._get_conn() as conn: conn.execute( INSERT INTO tasks (task_id, user_id, agent_name, raw_message, status, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?), (task_id, user_id, agent_name, raw_message, status, created_at, created_at), ) conn.commit() def update_status(self, task_id: str, status: str, result: str None): with self._get_conn() as conn: conn.execute( UPDATE tasks SET status ?, result ?, updated_at ? WHERE task_id ?, (status, result, datetime.now().isoformat(), task_id), ) conn.commit() def get_task(self, task_id: str): with self._get_conn() as conn: row conn.execute(SELECT * FROM tasks WHERE task_id ?, (task_id,)).fetchone() if row: return dict(row) return None任务状态建议使用以下几种枚举状态码含义下一步PENDING已经入队等待执行Worker 轮询取任务RUNNING正在执行中等待执行完成SUCCESS执行成功回传结果FAILED执行失败记录错误日志回传失败原因6.2 异步 WorkerWorker 的核心逻辑是从队列取出任务调用 Agent handler然后更新状态。如果是长时间执行的任务还可以在 Worker 里引入线程池。# 文件路径agent_scheduler/core/worker.py import queue import threading from typing import Dict from core.agent_registry import AgentSpec from core.result_sender import ResultSender from core.task_store import TaskStore class TaskWorker: def __init__(self): self._task_queue: queue.Queue queue.Queue() self._running True self._task_store TaskStore() self._result_sender ResultSender() self._thread threading.Thread(targetself._consume_loop, daemonTrue) self._thread.start() def submit(self, task_id: str, agent: AgentSpec, user_id: str, message: str): self._task_queue.put( { task_id: task_id, agent: agent, user_id: user_id, message: message, } ) def _consume_loop(self): while self._running: try: item self._task_queue.get(timeout1) except queue.Empty: continue task_id item[task_id] agent item[agent] user_id item[user_id] message item[message] self._task_store.update_status(task_id, RUNNING) try: # 调用 Agent 的执行函数 result agent.handler(user_id, message) self._task_store.update_status(task_id, SUCCESS, result) self._result_sender.send(user_id, f任务 {task_id} 执行完成\n{result}) except Exception as exc: error_msg str(exc) self._task_store.update_status(task_id, FAILED, error_msg) self._result_sender.send(user_id, f任务 {task_id} 执行失败{error_msg})Worker 使用了一个守护线程从队列里消费任务。这种设计在单机场景下已经够用而且不会因为某个 Agent 卡住导致系统崩溃。生产环境如果任务量大、Agent 执行时间长就应该把queue.Queue替换成 Redis 队列或者 RabbitMQ再把 Worker 拆成独立进程横向扩容。7. 结果回传把 Agent 执行结果推回聊天窗口任务执行完成之后结果如何回到用户聊天窗口这里要根据消息来源做不同处理。如果是企业微信或公众号通常用“主动推送消息”API如果是个人微信库直接通过 itchat 发送即可。# 文件路径agent_scheduler/core/result_sender.py import requests class ResultSender: 结果回调发送器。 以企业微信主动推送消息为例。 WECOM_API https://qyapi.weixin.qq.com/cgi-bin/message/send def __init__(self): self._access_token None self._token_expire_time 0 def _get_access_token(self): # 简化处理实际项目中应该使用 corpid secret 获取并缓存 token return your-access-token def send(self, user_id: str, content: str): # 企业微信文本消息格式 payload { touser: user_id, msgtype: text, agentid: 1000002, text: {content: content}, safe: 0, } headers {Content-Type: application/json} resp requests.post( f{self.WECOM_API}?access_token{self._get_access_token()}, jsonpayload, headersheaders, timeout10, ) resp.raise_for_status() return resp.json()这段代码在真实项目中需要注意两点access_token 的有效期是 7200 秒不能每次都获取必须做缓存。企业微信对主动推送消息的条数和频率有限制需要合理设计发送策略防止触发风控。如果你使用个人微信协议方式ResultSender.send可以改成调用 itchat 的send_msg接口逻辑更简单但稳定性也会弱一些。8. 完整示例注册一个能查询天气 / 执行命令的 Agent为了让整个系统连贯起来这里注册两个示例 Agent。一个用于演示调用外部 API天气查询另一个用于演示执行本地命令简单脚本。# 文件路径examples/register_agents.py import subprocess from agent_scheduler.core.agent_registry import registry, AgentSpec def weather_agent_handler(user_id: str, message: str) - str: 天气查询 Agent。 实际项目中可以调用和风天气、气象台 API 等。 # 从消息中解析城市名简化处理取符号后的全部内容作为城市 city message.split( , 1)[-1].strip() if in message else 北京 # 模拟请求外部 API真实实现请替换为 requests 调用 return f【天气查询】{city} 今日多云气温 18~27 摄氏度东南风 2 级。 def ping_agent_handler(user_id: str, message: str) - str: 演示型 Agent执行一段本地命令并返回结果。 注意生产环境必须严格限制命令执行权限建议使用白名单。 allowed_commands [echo, df -h] command message.split( , 1)[-1] if in message else echo hello if not any(cmd in command for cmd in allowed_commands): return 该命令未在白名单中已拒绝执行。 result subprocess.run(command, shellTrue, capture_outputTrue, textTrue, timeout10) return f命令执行结果\n{result.stdout} def register_demo_agents(): registry.register( AgentSpec( nameweather, description查询城市天气, keywords[天气, weather, 气温], handlerweather_agent_handler, ) ) registry.register( AgentSpec( nameshell, description执行白名单内的Shell命令, keywords[执行命令, shell, 运行], handlerping_agent_handler, ) )需要强调的是shell这个 Agent 只是为了演示调度链路不是鼓励大家直接暴露命令执行能力。如果生产环境确实需要执行命令一定要做到用户权限校验、命令白名单、超时控制、执行日志审计。8.1 启动入口最终的启动入口如下# 文件路径main.py from agent_scheduler.message_handlers.wecom import app from examples.register_agents import register_demo_agents if __name__ __main__: # 1. 注册所有 Agent register_demo_agents() # 2. 启动 Web 回调服务 # 这里以企业微信回调为例如果用个人微信就调用 start_wechat_bot() app.run(host0.0.0.0, port8000, debugFalse)9. 运行验证与效果检查把代码跑起来后怎么确认整个调度链路是通的建议按下面几步验证。9.1 基础验证Agent 注册状态启动服务后控制台会输出注册信息[AgentRegistry] 已注册 Agent: weather [AgentRegistry] 已注册 Agent: shell如果这一步输出缺失检查 import 路径和注册函数是否被调用。9.2 消息调度验证在企业微信中给自建应用发送消息weather 上海天气预期收到三类响应第一秒内任务已受理并给出任务号。稍后几秒Agent「weather」执行完成附带天气查询结果。任务状态表对应记录为 SUCCESS。如果直接查询 SQLite 表sqlite3 agent_tasks.db select * from tasks;能看到类似这样的记录task_id | user_id | agent_name | raw_message | status | result abc123 | zhangsan | weather | weather 上海天气 | SUCCESS | 【天气查询】上海...这就说明消息接入、意图匹配、异步执行、状态落库、结果回传这条链路完全跑通了。9.3 失败情况验证再故意发送一条触发 shell Agent 的消息但命令不在白名单内shell rm -rf /tmp/test预期返回结果中会包含“命令未在白名单中”的提示并且任务状态为 SUCCESS因为 handler 本身没有抛异常。这种“业务上拒绝但系统执行成功”的状态区分也很重要方便你后续在状态表基础上做更细粒度的权限审计。10. 常见问题与排查思路实际开发中几个高频问题基本集中在消息接入、签名校验和任务状态不一致上。问题现象可能原因排查方式解决方案企业微信回调 URL 验证失败token/signature 不匹配对比生成的签名和微信服务器传入的参数仔细核对 token 和 encodingAESKey按文档重新计算签名邮件/代码上消息收到了但调度没有响应消息入口没有调用dispatch_message在回调函数入口打印日志确认请求到达检查路由 URL 和请求方法Agent 执行成功但用户收不到结果access_token 过期或 API 频率限制查看 ResultSender 日志查看返回的错误码实现 token 缓存和过期刷新任务状态一直 PENDINGWorker 线程未启动查看 Worker 初始化日志确认TaskWorker()在进程启动时被实例化多个用户并发执行时互相干扰任务状态存储使用全局变量检查是否有连接共享导致的串数据使用 SQLite 行级写锁或迁移到 MySQL / Redis个人微信协议频繁掉线服务端对网页端协议限制无根治方案生产环境改用企业微信 / 公众号官方接口11. 工程落地时的最佳实践建议前面给出了一个可以跑通的最小系统但真正要长期稳定运行还需要引入下面这些工程约束。11.1 消息层安全与校验必须做消息来源合法性校验。企业微信回调要验证msg_signature个人微信方案要处理异常登录提示。不要裸奔的消息内容。如果有敏感信息建议在发送前做脱敏处理。对用户消息做长度限制防止超长消息导致 Agent 解析异常。11.2 调度层设计为可替换调度模块是系统扩展的核心。建议把dispatch_message设计成纯函数不直接依赖 Any 特定的微信库或 Web 框架。这样后续无论接钉钉、飞书还是 Slack都只需要新增一个消息适配器。# 伪代码示意消息适配器 class DingTalkAdapter: def handle(self, payload): user_id extract_user_id(payload) content extract_content(payload) return dispatch_message(user_id, content)11.3 状态层务必持久化不要用内存字典保存任务状态。进程重启后所有任务就丢了排查问题时你完全不知道之前发生了什么。SQLite 对中小项目已经够用等任务量上来再迁移到 MySQL 或 Redis。11.4 Agent 层超时与重试Agent 如果在执行中卡住可能会导致 Worker 线程堆积。建议每个 Agent handler 内部都设置超时控制比如调用大模型 API 时设置timeout参数。对于可重试任务可以在 Worker 中增加失败重试机制但要注意区分“可重试异常”和“不可重试异常”。11.5 日志层链路追踪建议为每个任务打印结构化的日志至少包含 task_id、user_id、agent_name、耗时。后续做性能和问题分析时这些都是基础数据。2025-01-01 12:00:01 INFO [dispatcher] 收到消息 task_idabc123 user_idzhangsan contentweather 上海天气 2025-01-01 12:00:03 INFO [worker] 任务开始执行 task_idabc123 agentweather 2025-01-01 12:00:04 INFO [worker] 任务执行成功 task_idabc123 cost_ms80012. 总结与后续扩展方向到目前为止一个“聊天窗口即调度台”的 Agent 系统已经完整落地了。从消息接入、意图匹配、任务调度、状态存储到结果回传每一层都有清晰代码可以验证。这条链路最难能可贵的部分在于它不依赖重量级中间件只需要一个 Python 环境加一个社交软件官方接口就能跑起来。后续如果想往更完整的工程化方向演进可以从三个方向继续深入。第一把单机 Worker 替换为独立的消息队列架构引入 Redis/RabbitMQ让 Agent 执行节点可以横向扩展。当你手上有大量 Agent且单个任务执行时间超过几分钟时这种架构是必要的。第二把关键词匹配替换为基于大模型的意图识别和函数调用。让调度中枢具备更好的语义理解能力例如用户说“今天适合穿什么”系统能自动匹配到天气 Agent 和穿衣建议 Agent并做结果聚合。这已经是目前主流的 Agent 编排方向。第三补充会话记忆。现在的dispatch_message是无状态的每次用户发消息都是独立任务。如果用户说出“后来呢”这类指代不明的话调度器无法理解。引入会话上下文缓存后Agent 之间的多轮协作才会真正自然起来。本文提供的这套“社交软件 Agent 调度”思路核心不在于某个代码有多精巧而在于它把 Agent 的入口问题解决掉了。当用户不再需要学习复杂的工具界面只需在聊天窗口里发一句话就能完成任务调度时Agent 的使用门槛已经被拉到和日常聊天一样的水平这在企业内部效率工具和自动化场景中有着非常实际的价值。建议你在自己的项目里先跑通最小闭环再逐步补充鉴权、多 Agent 编排和生产级消息队列。
返回列表