MemOS 调度器状态监控接口实战:任务进度、Redis 队列积压与系统概览全解析
【免费下载链接】MemOSSelf-evolving memory OS for LLM & AI Agents: ultra-persistent memory, hybrid-retrieval, and cross-task skill reuse, with 35.24% token savings and DeepSeek Harness support.项目地址: https://gitcode.com/gh_mirrors/memos/MemOS
导读
MemOS 的异步记忆生产链路(LLM 记忆提取、向量索引构建等)全部由 MemScheduler 调度体系在后台执行。为了让开发者实时掌握任务生命周期,MemOS 在开源版 Server(server_api)中提供了三个基于/product路由前缀的调度器状态监控接口:任务进度查询(/status)、用户队列指标(/task_queue_status)与系统级概览(/allstatus)。本文以 docs/cn/open_source/open_source_api/scheduler/get_status.md 为骨架,结合路由、Handler 与 Redis 队列的源码实现,完整讲解三个接口的参数、返回字段、底层工作原理与可落地的 Python 轮询示例,帮助你快速构建属于自己的记忆任务可观测系统。
1. 核心机理:MemScheduler 调度体系
在开源架构中,MemScheduler负责处理所有高耗时的后台任务(如 LLM 记忆提取、向量索引构建等)。理解状态监控接口之前,需要先掌握调度体系的三个基本事实:
- 状态流转:任务在生命周期内会经历
waiting(等待中)、in_progress(执行中)、completed(已完成)或failed(失败)等状态。从响应模型看,完整状态集合还包括pending与cancelled(见 product_models.py 中StatusResponseItem.status的Literal["in_progress", "completed", "waiting", "failed", "cancelled"]定义)。 - 队列监控:系统基于 Redis Stream 实现任务分发。通过监控
pending(已交付未确认)和remaining(排队中)任务数,可以评估系统的处理压力。 - 多维度观测:支持从"单任务"、"单用户队列"以及"全系统 summary"三个维度进行状态透视。
从调度器源码看,BaseScheduler在初始化时通过self.config.get("use_redis_queue", DEFAULT_USE_REDIS_QUEUE)决定使用 Redis 队列还是本地内存队列(见 base_scheduler.py);而环境变量开关MEMSCHEDULER_USE_REDIS_QUEUE默认为False,显式设置为"true"时才启用 Redis 队列(见 general_schemas.py)。Redis Stream 的键前缀固定为scheduler:messages:stream:v2.0(见 task_schemas.py)。
2. 接口详解
三个接口均由开源版 Server(server_api,路由前缀/product)直接提供,使用标准 HTTPGET请求即可访问。路由定义集中在 server_router.py 的 "Scheduler API Endpoints" 段落,Router 前缀在 server_router.py 处声明为/product。
2.1 任务进度查询(GET /product/scheduler/status)
用于追踪特定异步任务的当前执行阶段。对应路由实现见 server_router.py:
| 参数名 | 类型 | 必填 | 说明 |
|---|---|---|---|
user_id | str | 是 | 请求查询的用户唯一标识符。 |
task_id | str | 否 | 可选。若提供,则仅查询该特定任务的状态。 |
返回状态说明:
waiting: 任务已进入队列,等待空闲 Worker 执行。in_progress: Worker 正在调用大模型提取记忆或写入数据库。completed: 记忆已成功持久化并完成向量索引同步。failed: 任务失败。cancelled: 任务被取消(源码StatusResponseItem允许该取值)。
一个值得注意的细节:task_id存在两种语义(源码注释明确说明了这一点,见 scheduler_handler.py):
- business_task_id:业务级任务 ID,一个业务任务可拆分为多个 item 子任务,查询时会聚合所有关联 item 的状态;
- item_id:单个子任务的内部 ID,查询时返回单条状态。
Handler 会先尝试按 business_task_id 做聚合查询,失败后再按 item_id 查询;两者都查不到时返回 404("Task {task_id} not found for user {user_id}")。
2.2 用户队列指标(GET /product/scheduler/task_queue_status)
用于监控指定用户在 Redis 中的任务积压情况。对应路由实现见 server_router.py:
| 参数名 | 类型 | 必填 | 说明 |
|---|---|---|---|
user_id | str | 是 | 需查询队列状况的用户 ID。 |
核心指标项:
pending_tasks_count: 已分发给 Worker 但尚未收到确认(Ack)的任务数。remaining_tasks_count: 当前仍在队列中排队等待分配的任务总数。stream_keys: 匹配到的 Redis Stream 键名列表。
从响应模型TaskQueueData(见 product_models.py)可以看到更完整的字段:除上述三项外,还返回user_name、mem_cube_id、users_count(当前出现在队列 Stream 中的去重用户数),以及两个明细字段pending_tasks_detail、remaining_tasks_detail(按{stream_key}:{count}格式列出每个 Stream 的逐项计数)。此外,若调度器队列不可用或未连接 Redis,接口会分别返回 503 与 404 错误。
2.3 系统级概览(GET /product/scheduler/allstatus)
获取调度器的全局运行概况,通常用于管理员后台监控。对应路由实现见 server_router.py:
核心返回信息:
scheduler_summary: 包含系统当前的负载与健康状况。all_tasks_summary: 所有正在运行及排队任务的聚合统计。
两个 summary 使用同一个TaskSummary模型(见 product_models.py),包含waiting、in_progress、pending、completed、failed、cancelled六个维度计数及total总数,且每个字段默认值为 0。二者的差异在于:all_tasks_summary覆盖所有被 TaskStatusTracker 跟踪的任务(跨用户),而scheduler_summary只聚焦当前调度器管理下的任务,并叠加队列监控器的实时数据。
3. 响应模型速览
根据 product_models.py 的定义,三个接口的响应结构如下(BaseResponse统一携带message与data字段):
/status→StatusResponse,data为StatusResponseItem[],每个元素形如{"task_id": "...", "status": "in_progress"},默认 message 为 "Memory get status successfully"。/task_queue_status→TaskQueueResponse,data为TaskQueueData,默认 message 为 "Scheduler task queue status retrieved successfully"。/allstatus→AllStatusResponse,data为AllStatusResponseData{scheduler_summary, all_tasks_summary},默认 message 为 "Scheduler status summary retrieved successfully"。
4. 工作原理:SchedulerHandler 深度解析
所有三个接口的请求最终都由 scheduler_handler.py 中对应的三个 Handler 函数处理,整体工作流程遵循"缓存检索 → 队列确认 → 指标聚合"三步:
- 缓存检索:首先从 Redis 状态缓存中查找
task_id对应的实时进度。 - 队列确认:若查询队列指标,Handler 会调用 Redis 统计指令(如
XLEN、XPENDING)分析 Stream 状态。 - 指标聚合:对于全局状态请求,Handler 会汇总所有活跃节点的指标,生成系统级的 summary 数据。
4.1/status:状态缓存的读写链路
任务状态的持久化层是TaskStatusTracker(见 status_tracker.py),它把每个用户的任务元数据写入 Redis Hash,键格式为memos:task_meta:{user_id}。任务在不同阶段会调用不同的记录方法:
task_submitted():写入初始状态waiting,并记录task_type、mem_cube_id、submitted_at;若提供了business_task_id,还会向memos:task_items:{user_id}:{business_task_id}这个 Set 中SADD当前 item_id,从而建立"业务任务 → 子任务"的映射(TTL 均为 7 天)。task_started():将状态更新为in_progress并写入started_at。task_completed()/task_failed():写入终态completed/failed及对应时间戳、错误信息。
/status的查询链路(见 scheduler_handler.py)因此有两种模式:
- 不带
task_id:调用get_all_tasks_for_user()返回该用户全部任务的状态列表; - 带
task_id:先调用get_task_status_by_business_id()做聚合查询——聚合规则为"任一 item 失败则整任务failed;存在in_progress或waiting则整任务in_progress;全部completed才判定completed"(见 status_tracker.py),查不到时再回退为get_task_status()单条查询。
4.2/task_queue_status:用 XPENDING 与 XLEN 量化积压
handle_task_queue_status(见 scheduler_handler.py)的统计逻辑非常直观:
- 从
mem_scheduler.memos_message_queue解包得到底层队列,并校验 Redis 连接是否可用(未连接时尝试auto_initialize_redis()懒初始化,仍失败则返回 503)。 - 通过
queue_wrapper.get_stream_keys()获取全部 Stream 键,按":{user_id}:"子串过滤出属于当前用户的 Stream。Stream 键格式为{prefix}:{user_id}:{mem_cube_id}:{task_label},因此用户 ID 可通过拆分键段(倒数第 3 段)解析出来,用于统计users_count。 - 对每个用户 Stream:
- 调用
XPENDING(stream, consumer_group)取第一个返回值作为pending计数(已投递、未 Ack); - 调用
XLEN(stream)作为remaining计数(仍排队待分配)。
- 调用
- 聚合所有 Stream 得到
pending_tasks_count/remaining_tasks_count,并保留逐 Stream 明细。
消费组名称默认scheduler_group,由 redis_queue.py 的RedisTaskQueue构造函数参数指定;Worker 通过XREADGROUP拉取消息、XACK确认、XAUTOCLAIM回收超时消息(_ensure_consumer_group会在 Stream 缺失时用XGROUP CREATE ... MKSTREAM自动重建)。
4.3/allstatus:全系统指标聚合
handle_scheduler_allstatus(见 scheduler_handler.py)优先采用流式聚合以避免把全部任务载荷加载进内存:
- 用
SCAN遍历memos:task_meta:*键,再对每个 Hash 用HSCAN分批读取字段,反序列化后按status计数(默认只统计 24 小时内的任务,超过max_age_seconds=86400的陈旧记录会被跳过以降低噪音)。 - 若 Redis 不可用,则回退到
get_all_tasks_global()全量加载再聚合(Docker / 无 Redis 部署场景)。 - 最后用
task_schedule_monitor.get_tasks_status()的实时队列数据覆盖scheduler_summary:Redis 队列模式下遍历scheduler:前缀的逐 Stream 条目(running→in_progress,remaining→waiting,pending→pending);本地内存队列模式下直接采用顶层{"running", "remaining", "pending"}扁平汇总。
关于本地队列的这一点,有专门的回归测试佐证:tests/api/test_scheduler_handler_allstatus.py 记录了 issue #1395 的根因——旧实现只聚合scheduler:前缀的逐 Stream 条目,导致MEMSCHEDULER_USE_REDIS_QUEUE=false时 summary 被清零;修复后本地队列的running/remaining/pending能正确映射到in_progress/waiting/pending(测试断言见该文件 L63-L95)。
5. 快速上手示例
这些接口由开源版 Server(server_api,路由前缀/product)直接提供,使用标准 HTTP 请求即可访问。以下示例轮询任务状态直至完成:
import time import requests # 自部署 MemOS Server 的地址(如启用了鉴权,请自行补充 Authorization 请求头) base_url = "http://localhost:8000" # 1. 系统级概览:查看整个 MemOS 系统的运行健康度 resp = requests.get(f"{base_url}/product/scheduler/allstatus", timeout=10) resp.raise_for_status() global_res = resp.json() print(f"系统运行概况: {global_res['data']['scheduler_summary']}") # 2. 队列指标监控:检查特定用户的任务积压情况 resp = requests.get( f"{base_url}/product/scheduler/task_queue_status", params={"user_id": "dev_user_01"}, timeout=10, ) resp.raise_for_status() queue_res = resp.json() print(f"排队中任务数: {queue_res['data']['remaining_tasks_count']}") print(f"已下发未确认任务数: {queue_res['data']['pending_tasks_count']}") # 3. 任务进度追踪:轮询特定任务直至结束 task_id = "task_888999" active_states = {"waiting", "pending", "in_progress"} while True: resp = requests.get( f"{base_url}/product/scheduler/status", params={"user_id": "dev_user_01", "task_id": task_id}, timeout=10, ) resp.raise_for_status() items = resp.json().get("data", []) # data 为状态列表:[{"task_id": ..., "status": ...}] statuses = {item["status"] for item in items} print(f"任务 {task_id} 当前状态: {statuses or '空'}") if not statuses or statuses.isdisjoint(active_states): break time.sleep(2)代码要点解读:
allstatus的data.scheduler_summary是TaskSummary对象,可直接按waiting、in_progress、completed、failed、pending、cancelled、total键取数;task_queue_status的data中,除了两个核心计数,还可遍历stream_keys/pending_tasks_detail/remaining_tasks_detail做逐队列钻取;/status返回的data是状态列表而非单个对象,即使只传一个task_id,也需用item["status"]逐条取出;聚合查询(business_task_id)时会返回聚合后的单条状态;- 轮询退出条件建议同时覆盖
cancelled(终端状态)与"空列表"(任务记录过期被清理),上述isdisjoint(active_states)写法已兼容。
6. 配套能力:wait 与 SSE 流式监控
除三个状态查询接口外,路由表中还提供了两个面向"等待调度器空闲"场景的配套接口(见 server_router.py),可在自动化流水线中配合使用:
POST /product/scheduler/wait:阻塞轮询/status,直到指定用户的所有任务都处于终态(completed/failed/cancelled),或超过timeout_seconds(默认 120s,poll_interval默认 0.5s)超时返回。实现见 scheduler_handler.py。GET /product/scheduler/wait/stream:以 Server-Sent Events(SSE,text/event-stream)形式周期推送{"active_tasks", "status": "running"|"idle"|"timeout", ...}心跳帧,直到空闲或超时,适合在 Web 界面上展示实时进度。实现见 scheduler_handler.py。
例如,当你在完成一批记忆写入后需要确保下游依赖已经落盘,可以先调用wait阻塞等待,再继续后续逻辑,避免竞态。
7. 常见问题与运维建议
| 现象 | 原因与排查 |
|---|---|
/task_queue_status返回 503 "Scheduler queue is not available" | 调度器未初始化消息队列(memos_message_queue为空),确认MEMSCHEDULER_USE_REDIS_QUEUE=true且调度器已启动。 |
| 返回 503 "Scheduler queue not connected to Redis" | 队列存在但 Redis 连接不可用,检查 Redis 地址、端口与网络连通性。 |
| 返回 404 "No scheduler streams found for user ..." | 该用户在 Redis 中不存在任何 Stream 键,说明尚无任务入队,或任务已全部消费且 Stream 已清理。 |
/allstatus的scheduler_summary全为 0 | 若运行在无 Redis 的 Docker 环境(use_redis_queue=false),需确认使用包含本地队列修复的版本;可参考 test_scheduler_handler_allstatus.py 中描述的场景。 |
/status返回 404 "Task ... not found" | 任务 ID 不存在或记录已超过 7 天 TTL 被 Redis 清理(task_submitted等写入方法统一设置 7 天过期)。 |
运维建议:把pending_tasks_count与remaining_tasks_count接入监控告警(例如 pending 持续增长说明 Worker 消费异常,需要关注XAUTOCLAIM回收与消费组状态);利用allstatus的all_tasks_summary做容量规划与故障复盘;状态轮询时注意控制频率(如 2s 一次),避免对 Redis 产生不必要的压力。
8. 源码阅读指引
如需深入源码细节,建议按以下路径阅读:
- 路由层:server_router.py——三个查询接口与 wait / SSE 接口的定义、参数与响应模型绑定;
- Handler 层:scheduler_handler.py——
handle_scheduler_allstatus/handle_scheduler_status/handle_task_queue_status的完整实现; - 状态缓存层:status_tracker.py——Redis Hash / Set 结构、状态流转方法与 business_task_id 聚合逻辑;
- 队列实现层:redis_queue.py——消费组、
XREADGROUP/XACK/XAUTOCLAIM与 Stream 键管理; - 监控层:task_schedule_monitor.py——
get_tasks_status对 Redis / 本地队列两种后端的差异化输出; - 测试佐证:test_scheduler_handler_allstatus.py——本地队列模式下 summary 映射的回归测试;
- 数据模型:product_models.py——三个响应模型与
TaskSummary的字段定义。
掌握了以上链路,你就能基于这三个接口构建出从"单任务追踪"到"全系统健康度"的完整记忆调度观测面板。
【免费下载链接】MemOSSelf-evolving memory OS for LLM & AI Agents: ultra-persistent memory, hybrid-retrieval, and cross-task skill reuse, with 35.24% token savings and DeepSeek Harness support.项目地址: https://gitcode.com/gh_mirrors/memos/MemOS
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考