MemOS 调度器状态监控接口实战:任务进度、Redis 队列积压与系统概览全解析
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 在开源版 Serverserver_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. 接口详解三个接口均由开源版 Serverserver_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.pybusiness_task_id业务级任务 ID一个业务任务可拆分为多个 item 子任务查询时会聚合所有关联 item 的状态item_id单个子任务的内部 ID查询时返回单条状态。Handler 会先尝试按 business_task_id 做聚合查询失败后再按 item_id 查询两者都查不到时返回 404Task {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_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→StatusResponsedata为StatusResponseItem[]每个元素形如{task_id: ..., status: in_progress}默认 message 为 Memory get status successfully。/task_queue_status→TaskQueueResponsedata为TaskQueueData默认 message 为 Scheduler task queue status retrieved successfully。/allstatus→AllStatusResponsedata为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_seconds86400的陈旧记录会被跳过以降低噪音。若 Redis 不可用则回退到get_all_tasks_global()全量加载再聚合Docker / 无 Redis 部署场景。最后用task_schedule_monitor.get_tasks_status()的实时队列数据覆盖scheduler_summaryRedis 队列模式下遍历scheduler:前缀的逐 Stream 条目running→in_progressremaining→waitingpending→pending本地内存队列模式下直接采用顶层{running, remaining, pending}扁平汇总。关于本地队列的这一点有专门的回归测试佐证tests/api/test_scheduler_handler_allstatus.py 记录了 issue #1395 的根因——旧实现只聚合scheduler:前缀的逐 Stream 条目导致MEMSCHEDULER_USE_REDIS_QUEUEfalse时 summary 被清零修复后本地队列的running/remaining/pending能正确映射到in_progress/waiting/pending测试断言见该文件 L63-L95。5. 快速上手示例这些接口由开源版 Serverserver_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, timeout10) 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}, timeout10, ) 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}, timeout10, ) 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默认 120spoll_interval默认 0.5s超时返回。实现见 scheduler_handler.py。GET /product/scheduler/wait/stream以 Server-Sent EventsSSEtext/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_QUEUEtrue且调度器已启动。返回 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_queuefalse需确认使用包含本地队列修复的版本可参考 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),仅供参考