数据截至 (上游 commit 85532420387c)
MemScheduler:异步摄取与激活记忆刷新
30 秒导读: 前面几章讲了记忆"怎么抽、怎么组织、怎么检索"——但这些活很重,不能卡在用户等回答的主线程里。MemScheduler 就是那层后台调度:把重排、去重、固化这类任务丢进一个队列,由后台线程池异步消费、按"任务标签"分发给对应处理器,并周期性把明文记忆固化成 KV-Cache 激活记忆。它不发明新算法,它决定这些算法在什么时候、由谁、以什么顺序被调用。
1. 这是什么(零基础也能懂)
一句话定义: MemScheduler 是 MemOS 的异步任务调度层——一个"你把活交给我,我在后台按轻重缓急慢慢干"的后台工人系统。
它要解决什么问题。 想象你在和一个带记忆的 AI 聊天。每说一句话,系统理论上都该做一堆事:把这句话读成结构化记忆、更新工作记忆的排序、去重、必要时把常用记忆"预热"成缓存。如果这些全在你等回复的那几秒里同步做完,你会等到崩溃。
思路:把重活挪到后台。 于是 MemOS 把这些任务拆成消息,丢进一个队列,主线程"提交完就返回",真正的计算交给后台线程池。这正是经典的生产者-消费者模式(一边往队列塞任务,一边有工人从队列取任务干)。
给谁用。 它不直接面向终端用户,而是被 MOSCore(第 2 章的内核)在内部挂接和驱动。开发者通过 enable_mem_scheduler 开关决定要不要启用它。
用起来什么样。 从内核视角,启用后只是多了一个后台服务在转:
# 示意,非源码:MOSCore 内部如何用调度器
scheduler = SchedulerFactory.from_config(scheduler_config) # 按 backend 造实例
scheduler.initialize_modules(chat_llm, process_llm, db_engine) # 装配子模块
scheduler.start() # 启动后台消费线程 + 监控线程
# ……此后主线程只管把消息 submit 进去,后台自己消费
scheduler.stop() # 优雅停机
一句话直觉: 把它想成餐厅后厨的叫号+传菜系统。前台(主线程)接单就走,订单进"票夹"(队列),后厨(线程池)按菜品类型(任务标签)分给不同灶台(处理器),还有个领班(监控)盯着谁卡住了。
2. 顶层全景(它大概怎么转)
先看一张"一条任务从提交到执行"的总图。怎么读: 从左到右是数据流;上半是"入队",下半是"出队消费"。
┌─────────────────────────────────────────────┐
MOSCore 调用 │ BaseScheduler │
submit_messages │ (骨架:装配子模块 + 管理 mem_cube + 生命周期) │
─────────────► │ │
└───────┬─────────────────────────────┬───────┘
│ 高优先级(LEVEL_1) │ 普通任务
│ 立即同步执行 ▼
│ ┌──────────────────┐
│ │ ScheduleTaskQueue │
│ │ 队列(二选一): │
│ │ · 本地内存队列 │
│ │ · Redis Streams │
│ └────────┬─────────┘
│ │
┌─────────────▼──────────────┐ │ 后台消费线程
│ SchedulerDispatcher │ ◄────────────┘ _message_consumer
│ 按 (user, cube, label) 分组 │ 循环:取一批 → dispatch
│ → 查 handler → 线程池执行 │
└─────────────┬──────────────┘
│ 按标签路由
┌─────────────────┼─────────────────┬──────────────┐
▼ ▼ ▼ ▼
query_handler add_handler mem_update_handler ……
(重排工作记忆) (写入记忆) (刷新激活记忆)
各部件一句话职责:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
BaseScheduler | 调度器骨架:装配子模块、管理 mem_cube、控制启停生命周期 | base_scheduler.py:69 |
GeneralScheduler | 具体实现:注册各任务标签的处理器 | general_scheduler.py:16 |
OptimizedScheduler | 增强实现:加 API 混合检索 + 更好的工作记忆替换 | optimized_scheduler.py:37 |
SchedulerFactory | 工厂:按配置 backend 造出对应调度器 | scheduler_factory.py:9 |
ScheduleTaskQueue | 队列封装:本地内存队列 / Redis Streams 二选一 | task_schedule_modules/task_queue.py:23 |
SchedulerDispatcher | 分发器:分组、查处理器、丢线程池执行 | task_schedule_modules/dispatcher.py:38 |
ActivationMemoryManager | 把明文记忆固化成 KV-Cache 激活记忆 | memory_manage_modules/activation_memory_manager.py:18 |
| 各监控模块 | 盯线程池健康、指标、任务状态 | monitors/ |
主线走一遍(高层): MOSCore 把一句话包成 ScheduleMessageItem → submit_messages 入队 → 后台 _message_consumer 循环取一批 → Dispatcher 按标签路由到 handler → handler 里可能触发"周期性固化激活记忆"。
3. 调度器骨架与生命周期
本节讲
BaseScheduler这个"底座":它怎么把一堆子模块装配起来,怎么管理记忆立方体(mem_cube),以及启停时都做了什么。
3.1 装配子模块:initialize_modules
调度器是个"空壳"直到它被喂进两个 LLM、一个数据库引擎。initialize_modules 一次性把监控、检索器、后处理器、激活记忆管理器全建起来:
# base_scheduler.py:225-242(节选)
self.monitor = SchedulerGeneralMonitor(process_llm=..., config=..., db_engine=...)
self.dispatcher_monitor = SchedulerDispatcherMonitor(config=self.config)
self.retriever = SchedulerRetriever(process_llm=self.process_llm, config=self.config)
self.post_processor = MemoryPostProcessor(process_llm=self.process_llm, config=self.config)
self.activation_memory_manager = ActivationMemoryManager(...)
这段真实实现见 base_scheduler.py:202 的 initialize_modules。关键设计: 若 enable_parallel_dispatch 为真,这里会顺手把 dispatcher_monitor 启动起来盯线程池(base_scheduler.py:247-249)。
3.2 失败即回滚:_cleanup_on_init_failure
装配是"要么全成、要么别留半拉子"。整段 initialize_modules 包在 try/except 里,一旦任何子模块建到一半抛错,就调 _cleanup_on_init_failure 把已启动的监控线程停掉,再把异常重新抛出:
# base_scheduler.py:270-274
except Exception as e:
logger.error(f"Failed to initialize scheduler modules: {e}", exc_info=True)
self._cleanup_on_init_failure() # 关掉已启动的 dispatcher_monitor
raise
_cleanup_on_init_failure 本体在 base_scheduler.py:276——目前只负责停 dispatcher_monitor,防止"初始化失败却留了个后台线程在转"这种资源泄漏。