MemOS 调度器状态监控接口实战:任务进度、Redis 队列积压与系统概览全解析
2026/9/23 11:58:22 网站建设 项目流程

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(失败)等状态。从响应模型看,完整状态集合还包括pendingcancelled(见 product_models.py 中StatusResponseItem.statusLiteral["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_idstr请求查询的用户唯一标识符。
task_idstr可选。若提供,则仅查询该特定任务的状态。

返回状态说明

  • 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_idstr需查询队列状况的用户 ID。

核心指标项

  • pending_tasks_count: 已分发给 Worker 但尚未收到确认(Ack)的任务数。
  • remaining_tasks_count: 当前仍在队列中排队等待分配的任务总数。
  • stream_keys: 匹配到的 Redis Stream 键名列表。

从响应模型TaskQueueData(见 product_models.py)可以看到更完整的字段:除上述三项外,还返回user_namemem_cube_idusers_count(当前出现在队列 Stream 中的去重用户数),以及两个明细字段pending_tasks_detailremaining_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),包含waitingin_progresspendingcompletedfailedcancelled六个维度计数及total总数,且每个字段默认值为 0。二者的差异在于:all_tasks_summary覆盖所有被 TaskStatusTracker 跟踪的任务(跨用户),而scheduler_summary只聚焦当前调度器管理下的任务,并叠加队列监控器的实时数据。


3. 响应模型速览

根据 product_models.py 的定义,三个接口的响应结构如下(BaseResponse统一携带messagedata字段):

  • /statusStatusResponsedataStatusResponseItem[],每个元素形如{"task_id": "...", "status": "in_progress"},默认 message 为 "Memory get status successfully"。
  • /task_queue_statusTaskQueueResponsedataTaskQueueData,默认 message 为 "Scheduler task queue status retrieved successfully"。
  • /allstatusAllStatusResponsedataAllStatusResponseData{scheduler_summary, all_tasks_summary},默认 message 为 "Scheduler status summary retrieved successfully"。

4. 工作原理:SchedulerHandler 深度解析

所有三个接口的请求最终都由 scheduler_handler.py 中对应的三个 Handler 函数处理,整体工作流程遵循"缓存检索 → 队列确认 → 指标聚合"三步:

  1. 缓存检索:首先从 Redis 状态缓存中查找task_id对应的实时进度。
  2. 队列确认:若查询队列指标,Handler 会调用 Redis 统计指令(如XLENXPENDING)分析 Stream 状态。
  3. 指标聚合:对于全局状态请求,Handler 会汇总所有活跃节点的指标,生成系统级的 summary 数据。

4.1/status:状态缓存的读写链路

任务状态的持久化层是TaskStatusTracker(见 status_tracker.py),它把每个用户的任务元数据写入 Redis Hash,键格式为memos:task_meta:{user_id}。任务在不同阶段会调用不同的记录方法:

  • task_submitted():写入初始状态waiting,并记录task_typemem_cube_idsubmitted_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_progresswaiting则整任务in_progress;全部completed才判定completed"(见 status_tracker.py),查不到时再回退为get_task_status()单条查询。

4.2/task_queue_status:用 XPENDING 与 XLEN 量化积压

handle_task_queue_status(见 scheduler_handler.py)的统计逻辑非常直观:

  1. mem_scheduler.memos_message_queue解包得到底层队列,并校验 Redis 连接是否可用(未连接时尝试auto_initialize_redis()懒初始化,仍失败则返回 503)。
  2. 通过queue_wrapper.get_stream_keys()获取全部 Stream 键,按":{user_id}:"子串过滤出属于当前用户的 Stream。Stream 键格式为{prefix}:{user_id}:{mem_cube_id}:{task_label},因此用户 ID 可通过拆分键段(倒数第 3 段)解析出来,用于统计users_count
  3. 对每个用户 Stream:
    • 调用XPENDING(stream, consumer_group)取第一个返回值作为pending计数(已投递、未 Ack);
    • 调用XLEN(stream)作为remaining计数(仍排队待分配)。
  4. 聚合所有 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)优先采用流式聚合以避免把全部任务载荷加载进内存:

  1. SCAN遍历memos:task_meta:*键,再对每个 Hash 用HSCAN分批读取字段,反序列化后按status计数(默认只统计 24 小时内的任务,超过max_age_seconds=86400的陈旧记录会被跳过以降低噪音)。
  2. 若 Redis 不可用,则回退到get_all_tasks_global()全量加载再聚合(Docker / 无 Redis 部署场景)。
  3. 最后用task_schedule_monitor.get_tasks_status()的实时队列数据覆盖scheduler_summary:Redis 队列模式下遍历scheduler:前缀的逐 Stream 条目(runningin_progressremainingwaitingpendingpending);本地内存队列模式下直接采用顶层{"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)

代码要点解读

  • allstatusdata.scheduler_summaryTaskSummary对象,可直接按waitingin_progresscompletedfailedpendingcancelledtotal键取数;
  • task_queue_statusdata中,除了两个核心计数,还可遍历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 已清理。
/allstatusscheduler_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_countremaining_tasks_count接入监控告警(例如 pending 持续增长说明 Worker 消费异常,需要关注XAUTOCLAIM回收与消费组状态);利用allstatusall_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),仅供参考

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

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

立即咨询