MathModelAgent如何实现任务队列?Redis发布订阅机制的深度解析
【免费下载链接】MathModelAgent🤖📐专为数学建模设计的 Agent & skills ,自动完成数学建模,生成一份完整的可以直接提交的论文。 An Agent Designed for Mathematical Modeling ,Automatically complete mathmodel and generate a complete paper ready for submission.项目地址: https://gitcode.com/GitHub_Trending/ma/MathModelAgent
MathModelAgent 是一款专为数学建模场景设计的 Agent 项目,它能自动完成选题分析、建模、编码求解到论文写作的全流程。由于一个建模任务往往要运行数小时,项目必须解决两个核心问题:任务如何排队执行?执行进度如何实时推送给前端?答案正是本文要拆解的Redis 发布订阅(Pub/Sub)+ WebSocket 任务队列机制。
一、整体架构:一次"提交即返回"的异步任务
传统 Web 服务里,用户提交一个耗时任务后,HTTP 连接要一直挂着等待结果——这对动辄数小时的建模任务完全不可行。MathModelAgent 采用的是经典的异步任务队列模式:
- 入队:用户提交任务,后端立即生成
task_id并创建独立工作目录; - 快速响应:接口马上返回
{"task_id": ..., "status": "processing"},HTTP 连接即刻释放; - 后台执行:真正的建模流程被丢进后台,通过 Redis 频道持续"广播"进度;
- 实时消费:前端通过 WebSocket 订阅该任务的频道,逐条接收进度消息。
入口逻辑位于 modeling_router.py,提交接口的关键动作只有三步:
await redis_manager.set(f"task_id:{task_id}", task_id) # 任务注册到 Redis background_tasks.add_task(run_modeling_task_async, ...) # 任务入队后台执行 return {"task_id": task_id, "status": "processing"} # 立即返回这里有个巧妙的设计:task_id:{task_id}这个键写入 Redis 后设置了36000 秒(10 小时)过期时间,它既是"任务是否存在"的凭证,也天然自动清理僵尸任务,无需额外维护。
二、任务如何执行:asyncio 任务注册表
任务真正跑起来的地方是run_modeling_task_async。它做了三件关键的事:
- 发布"任务开始"消息到 Redis 频道,让前端立刻看到状态变化;
asyncio.create_task启动建模工作流,并设置 5 小时超时(asyncio.wait_for);- 将
(asyncio.Task, asyncio.Event)注册进全局字典_active_tasks——这正是实现"一键停止任务"的基础。
任务结束时,无论成功、取消还是异常,都会再发布一条对应的状态消息(success/warning/error),保证前端状态永远有收尾。
三、核心机制:Redis 发布订阅频道设计
消息发布的核心封装在 redis_manager.py 中,每个任务对应一个独立频道:
频道命名规则:task:{task_id}:messagespublish_message方法做了双写——既PUBLISH到 Redis 频道实现实时推送,又把消息追加保存到logs/messages/{task_id}.json文件做持久化。为什么需要持久化?因为 Redis Pub/Sub 的消息不落地,如果用户刷新页面时 WebSocket 还没连上,中间的消息就丢了。文件兜底后,前端随时可以通过 common_router.py 的GET /messages?task_id=...接口拉取完整历史消息,实现"断点续看"。
消息体本身由 response.py 中一组 Pydantic 模型定义,按msg_type区分为 system(系统状态)、agent(四个 Agent 的输出)、tool(代码执行结果)等类型,前端据此渲染不同样式的气泡与 Notebook 单元格。
四、WebSocket 订阅端:从 Redis 到浏览器的最后一跳
ws_router.py 定义了WebSocket /task/{task_id}端点,它的工作流程是:
- 校验任务:检查
task_id:{task_id}键是否存在,不存在直接关闭连接(防止访问非法任务); - 订阅频道:调用
redis_manager.subscribe_to_task()拿到该任务的 Pub/Sub 订阅句柄; - 消息转发循环:以 0.1 秒为节拍轮询
pubsub.get_message(),解析 JSON 后原样推送给浏览器; - 优雅退出:连接关闭时取消订阅并移除连接记录(由 ws_manager.py 统一管理活跃连接)。
还有一个体现细节的小设计:任务启动后会有await asyncio.sleep(1)的短暂延迟——这是为了给前端 WebSocket 留出连接窗口,避免"开始处理"这条首条消息因没人订阅而丢失(尽管有文件兜底,实时性上依然尽力保证)。
五、前端配套:指数退避自动重连
浏览器侧的 websocket.ts 实现了TaskWebSocket类,它不只是简单建连,还内置了一套生产级的重连策略:
- 断线后自动重连,最多10 次;
- 采用指数退避:间隔从 1 秒翻倍递增,上限 30 秒,避免服务端故障时疯狂重试打爆服务;
- 手动关闭(
close())时置位isManualClose标记,跳过自动重连逻辑。
配合后端的文件持久化 +/messages历史接口,即使长任务中途断网重连,用户刷新页面也能恢复完整进度,这正是任务队列体验的"最后一公里"。
六、一键停止任务:取消信号的传递
建模任务跑偏了怎么办?项目通过_active_tasks注册表 +asyncio.Event实现协作式取消:
_, cancel_event = _active_tasks[task_id] cancel_event.set() # 仅发信号,工作流内部周期性检查并自行优雅退出POST /modeling/{task_id}/cancel接口只负责"发信号",工作流执行到下一个检查点时感知到事件后自行收尾,从而避免粗暴杀进程导致的资源泄漏。
七、部署细节:一条 Docker Compose 拉起 Redis
在 docker-compose.yml 中,Redis 以redis:alpine镜像运行并挂载持久化卷,后端通过depends_on保证 Redis 先行启动。连接地址由 setting.py 中的两个配置项控制:
REDIS_URL:默认redis://redis:6379/0(Docker 网络内服务发现);REDIS_MAX_CONNECTIONS:连接池上限,默认 10,适配单机中小并发场景。
此外GET /status健康检查接口会实时 ping Redis 并返回后端与缓存服务的运行状态,方便运维观测。
总结:这套任务队列方案值得借鉴的 4 个点
| 设计点 | 实现方式 | 解决的问题 |
|---|---|---|
| 异步任务入队 | FastAPI BackgroundTasks + asyncio.Task | 长任务不阻塞 HTTP |
| 进度广播 | Redis Pub/Sub 按任务隔离频道 | 多任务并发、互不干扰 |
| 消息双写 | Pub/Sub 实时推送 + JSON 文件持久化 | 断线不丢进度、可回溯 |
| 优雅取消 | asyncio.Event 协作式信号 | 随时停止且资源安全回收 |
整体来看,MathModelAgent 用"Redis 发布订阅 + WebSocket + 文件持久化"三件套,以极低的架构复杂度支撑起了多任务并发、实时进度、断线恢复与任务取消四大能力,是学习异步任务队列设计的一个非常贴切、可读的开源范例。相关源码集中在 backend/app/services/ 与 backend/app/routers/ 目录,配合 docs/tutorial.md 教程可以快速上手本地部署。
【免费下载链接】MathModelAgent🤖📐专为数学建模设计的 Agent & skills ,自动完成数学建模,生成一份完整的可以直接提交的论文。 An Agent Designed for Mathematical Modeling ,Automatically complete mathmodel and generate a complete paper ready for submission.项目地址: https://gitcode.com/GitHub_Trending/ma/MathModelAgent
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考