MathModelAgent如何实现任务队列?Redis发布订阅机制的深度解析
2026/9/17 14:21:47 网站建设 项目流程

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 采用的是经典的异步任务队列模式:

  1. 入队:用户提交任务,后端立即生成task_id并创建独立工作目录;
  2. 快速响应:接口马上返回{"task_id": ..., "status": "processing"},HTTP 连接即刻释放;
  3. 后台执行:真正的建模流程被丢进后台,通过 Redis 频道持续"广播"进度;
  4. 实时消费:前端通过 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}:messages

publish_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}端点,它的工作流程是:

  1. 校验任务:检查task_id:{task_id}键是否存在,不存在直接关闭连接(防止访问非法任务);
  2. 订阅频道:调用redis_manager.subscribe_to_task()拿到该任务的 Pub/Sub 订阅句柄;
  3. 消息转发循环:以 0.1 秒为节拍轮询pubsub.get_message(),解析 JSON 后原样推送给浏览器;
  4. 优雅退出:连接关闭时取消订阅并移除连接记录(由 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询