跳到主要内容

图谱式明文记忆(上):组织、去重、冲突与重构

30 秒导读: MemReader(见 03)把对话/文档嚼成一条条结构化"记忆项"之后,这些碎片要落进一张图长期维护——这就是本章。你会看到:一条记忆项如何被并发写成图节点、系统怎么在后台悄悄判断"这条和已有的重复/矛盾了"并合并、旧版本如何被归档而非硬删、以及一个定时线程如何把散落的节点聚类、生成"摘要父节点",把平铺的碎片重构成一棵树。检索(怎么把它们再捞出来)是下一章 05,本章只讲"写入图谱"这一侧。


1. 这一章讲什么(先划清边界)

MemOS 的明文记忆(TreeTextMemory)本质是一张存在图数据库里的知识图:节点是记忆,边是记忆之间的关系(父子、合并、因果……)。

围绕这张图有两件事要做,本章只做第一件:

做什么在哪讲
写入 / 组织(本章)记忆项 → 图节点;去重、冲突消解、归档、后台层级重构本章
检索 / 召回查询 → 命中节点 → 重排 → 推理 → 返回05-retrieval-pipeline.md

一句话记住定位:本章是"往图里写、并把图养好"的那一半。

一条记忆项(TextualMemoryItem)是什么样,回顾 01;这里只需要知道它有三个关键部分:memory(正文)、metadata(元数据,含 embeddingmemory_typestatustags 等)、id


2. 顶层全景:一条记忆从"项"到"图里的节点"

先看大盘。从 TreeTextMemory.add() 进来,到最终落图,链路分**同步的"快写"异步的"慢养"**两段:

TreeTextMemory.add(memories) tree.py:103
│ (只是转发)

MemoryManager.add(...) manager.py:89
┌──────────────┴───────────────┐
use_batch=True use_batch=False
_add_memories_batch _add_memories_parallel
manager.py:138 manager.py:120
│ │ (每条一个线程)
│ ▼
│ _process_memory manager.py:309
└──────────┬─────────────────┘

①同步:把节点写进图数据库(graph_store.add_node)
│ —— 这一步返回,调用方就拿到 id 了

②异步:塞一条 QueueMessage(op="add") 进重构队列
│ manager.py:243 / 420
═══════════════════════╪═══════════ 线程边界(后台) ═══════════

GraphStructureReorganizer reorganizer.py:81
┌────────────────┴─────────────────┐
消息消费线程 定时优化线程
_run_message_consumer_loop _run_structure_organizer_loop
reorganizer.py:133 reorganizer.py:151 (每 100s)
│ │
▼ ▼
handle_add → NodeHandler optimize_structure
去重/冲突消解 _partition → _summarize_cluster
handler.py:30/76 聚类 + 生成摘要父节点

怎么读这张图:从上到下是时间顺序;中间那条虚线是"同步返回点"add() 在①之后就返回了——调用方立刻拿到新节点 id,不必等去重和重构。②之后的所有"整理"工作(去重、消解冲突、聚类成树)都在后台线程里慢慢做,不阻塞写入

各部件一句话职责:

部件干什么文件
TreeTextMemory明文记忆的门面,add/search/soft_delete/drop 等对外方法tree.py:39
MemoryManager写入编排:决定批量还是并行、双写工作记忆+图节点、维护容量organize/manager.py:54
GraphStructureReorganizer后台守护:消费队列做增量整理 + 定时做全局聚类重构organize/reorganizer.py:81
NodeHandler增量整理的核心:检测"重复/矛盾"并消解(合并或硬更新)organize/handler.py:22
RelationAndReasoningDetector聚类时的关系/推理探测(因果、条件、推断节点)organize/relation_reason_detector.py:20
BaseGraphDB + 后端图数据库抽象层 + Neo4j/Postgres/PolarDB 实现graph_dbs/base.py:11

3. 写入编排:一条记忆怎么落图

本节讲 MemoryManager——它是"写"这一侧的调度中枢。

3.1 两条写入路径:批量 vs 并行

MemoryManager.add() 是入口,它按 use_batch 分岔(默认 True):

# 示意,非源码 —— manager.py:89 add() 的骨架
def add(self, memories, user_name=None, mode="sync", use_batch=True):
if use_batch:
added_ids = self._add_memories_batch(memories, user_name) # 批量:少次数据库往返
else:
added_ids = self._add_memories_parallel(memories, user_name) # 并行:每条一个线程
if mode == "sync":
self._cleanup_working_memory(user_name) # 同步模式顺手清理超量的工作记忆
return added_ids
  • 并行路径 _add_memories_parallel(manager.py:120):开一个 10 线程的池,每条记忆丢给 _process_memory 单独处理。适合零散写入。
  • 批量路径 _add_memories_batch(manager.py:138):先把所有记忆整理成节点字典列表,再按 batch_size=5 分批调 graph_store.add_nodes_batch。适合一次灌大量记忆时省数据库往返。

mode("sync"/"async")来自 TreeTextMemory 的配置(tree.py:114self.mode 透传进来),决定写完要不要顺手做工作记忆清理。

3.2 双写:一条记忆既是"工作记忆"也是"长期记忆"

这是最容易看漏的设计。_process_memory(manager.py:309)对一条 LongTermMemory/UserMemory并行写两份:

  • 一份 WorkingMemory 节点(短时、FIFO 淘汰)——_add_memory_to_db(..., "WorkingMemory"),manager.py:369
  • 一份对应类型的图节点(长期)——_add_to_graph_memory,manager.py:390
# 示意,非源码 —— _process_memory 的核心(manager.py:329)
with ContextThreadPoolExecutor(max_workers=2) as ex:
if memory_type in ("WorkingMemory","LongTermMemory","UserMemory","OuterMemory"):
ex.submit(self._add_memory_to_db, memory, "WorkingMemory", ...) # 写工作记忆
if memory_type in ("LongTermMemory","UserMemory","RawFileMemory", ...):
ex.submit(self._add_to_graph_memory, memory=memory, ...) # 写长期图节点
# 只返回长期节点的 id,不返回工作记忆 id —— 保持对外契约

两份靠一个 working_binding 字段互相记账:长期节点的 metadata 里记下它对应的工作记忆 id(manager.py:404)。在"快模式"(mode:fast 标签)下,还会把 [working_binding:<uuid>] 这行塞进 background,以便日后 MemReader 产出精修版记忆后,能反查并清掉那批临时的 WorkingMemory 节点(见 extract_working_binding_ids,manager.py:23)。

一个坑(诚实提醒): 批量路径 _add_memories_batch 里,写 WorkingMemory 那段被注释停用了(manager.py:239-241 的 TODO:"working id is same with item.id, need to fix"),所以批量写目前只落长期图节点、不落工作记忆。只有并行路径 _process_memory 才真正双写。两条路径行为并不完全对齐。

3.3 写完就排队:交给后台整理

无论哪条路径,写进图之后都会往重构队列塞一条 QueueMessage(op="add", after_node=[...]):

  • 并行路径:_add_to_graph_memory 每写一个节点就排一条消息(manager.py:420)。
  • 批量路径:一整批只排一条消息,after_node 是这批所有节点 id(manager.py:243)。
# 示意,非源码 —— _add_to_graph_memory 尾部(manager.py:414)
self.graph_store.add_node(node_id, memory.memory, metadata_dict, user_name=user_name)
self.reorganizer.add_message( # 立刻返回,整理留给后台
QueueMessage(op="add", after_node=[node_id], user_name=user_name)
)
return node_id

第二个坑: 后台消费 add 消息时,handle_add 只取 after_node[0] 做去重(reorganizer.py:194)。这意味着批量路径里一条消息带 N 个节点,只有第 0 个会被拿去和存量做重复/冲突检测,其余 N-1 个这一轮被跳过(要等下一节的定时聚类才可能被处理到)。

3.4 两个"结构工具":ensure_structure_path 与 inherit_edges

MemoryManager 还带两个图结构工具,理解概念即可:

  • _ensure_structure_path(manager.py:459):保证一条结构路径存在(ROOT → … → 目标节点)。它按 metadata.key 查有没有同名结构节点,没有就新建一个作为"分类锚点"并返回其 id,给记忆当 parent。
  • _inherit_edges(manager.py:429):把一个节点的所有非血缘边(即除 MERGED_TO 之外的边)迁移到另一个节点上,迁移后删掉原边——合并节点时用来"继承邻居关系"。

诚实说明: 在本 commit 里,这两个方法已定义但未被 add 主链路调用(仓库内只有定义处引用),更像是为结构化写入预留的脚手架;真正在跑的"继承边"逻辑在冲突消解里(见 3.5 的 _resolve_in_graph)。


4. 关系推理与去重/冲突消解:让图不长满重复和矛盾

本节讲后台消息消费线程这条增量整理链:每加一个节点,就问一句"它和图里已有的谁重复了、和谁矛盾了?"并当场处理。核心是 NodeHandler(handler.py:22)。

4.1 检测:先向量召回,再让 LLM 判定关系

NodeHandler.detect(handler.py:30)分三步:

  1. 用新节点的 embedding 做向量搜索,取相似度 ≥ EMBEDDING_THRESHOLD(0.8,handler.py:23)的候选。
  2. 排除自身。
  3. 对每个候选,用 LLM(MEMORY_RELATION_DETECTOR_PROMPT)判定二者关系,只认三种输出:
LLM 判定含义后续动作
contradictory矛盾(如"他住北京"vs"他住上海")进入消解,尝试融合或按时间硬更新
redundant冗余重复进入消解,融合成一条
independent无关不处理
# 示意,非源码 —— detect 的判定循环(handler.py:49)
for cand in embedding_candidates:
result = llm.generate(DETECTOR_PROMPT.format(s1=memory.memory, s2=cand.memory)).strip()
if result == "contradictory":
detected.append([memory, cand, "contradictory"])
elif result == "redundant":
detected.append([memory, cand, "redundant"])
# independent / 其它:跳过

判定出来的每一对,交给 resolve 逐对消解(由 handle_add 驱动,reorganizer.py:200)。

4.2 消解:能融合就融合,融不了就按时间硬更新

NodeHandler.resolve(handler.py:76)让 LLM(MEMORY_RELATION_RESOLVER_PROMPT)读两条记忆 + 关键元数据,产出一段 <answer>…</answer>:

resolve(a, b, relation)
│ LLM 融合

<answer> 内容是什么? </answer>
│ │
很短且含 "no" 一段融合后的新记忆
(融不了) │
▼ ▼
_hard_update _resolve_in_graph
比 updated_at, 新建 merged 节点,
删旧留新 继承双方的边,
handler.py:131 把 a、b 标 archived
并连 MERGED_TO
handler.py:151

两条分支的关键差异:

  • 硬更新 _hard_update(handler.py:131):比较两者 updated_at,删掉旧的、保留新的(graph_store.delete_node 硬删)。用于 LLM 也无法调和的矛盾。
  • 图内融合 _resolve_in_graph(handler.py:151):
    1. 收集 a、b 的所有边;
    2. 新建 merged 节点(融合后的正文 + 合并元数据);
    3. 把两者的边改指到 merged(已存在则跳过);
    4. 把 a、b 的 status 置为 archived(不删,handler.py:184-185);
    5. 加两条 a -MERGED_TO-> mergedb -MERGED_TO-> merged 血缘边(handler.py:186-187)。

这样"旧版本"仍在图里(archived),血缘可追溯,而检索默认只看 activated。融合元数据由 _merge_metadata(handler.py:192)完成:sources 取并集、embedding 用融合正文重算、其余字段"谁非空取谁"。

4.3 另一条(当前休眠的)关系推理线

RelationAndReasoningDetector(relation_reason_detector.py:20)本意是更丰富的关系挖掘:成对判定 CAUSE/CONDITION/RELATE/CONFLICT(_parse_relation_result 的合法集,relation_reason_detector.py:227)、从因果对推断出新的"推理事实"节点(_infer_fact_nodes_from_relations,:118)、时序 FOLLOWS 边、聚合概念节点。

诚实提醒(重要): 在本 commit 里,process_node(:26)的四段实现(pairwise/inferred/sequence/aggregate)全部被三引号字符串注释掉了(:49-80),只保留了"若节点是 reasoning 类型就跳过"这一句。也就是说,这条更花哨的关系推理线目前实际不产出任何边;真正在跑的去重/冲突逻辑是 4.1–4.2 的 NodeHandler。聚类流程(下一节)虽然仍会调用它,但结果字典恒为空。


5. 归档与软删除:记忆怎么"退场"

MemOS 对"删除"很克制——尽量软删、保留版本,而不是抹掉。

5.1 状态机:一个节点的一生

节点 status 有四态(item.py:109):

activated ──冲突/重复被融合──▶ archived ──(可继续保留/清理)
│ ▲
│ │ MERGED_TO 边指向融合后的新节点
└──── soft_delete ────▶ deleted (打标,不真删)
  • activated:正常可被检索。
  • archived:被融合替代的旧版本,靠 MERGED_TO 挂在新节点后(见 4.2)。
  • deleted:软删标记。TreeTextMemory.soft_delete(tree.py:621)只是并发地把节点 status 改成 "deleted",可选再记一个 evolve_to(指向它"演化成"的新节点 id)——数据仍在。
  • 真正的物理删除只有 delete/_hard_updategraph_store.delete_node(tree.py:402handler.py:146)。

5.2 版本机制:ArchivedTextualMemory

当一条记忆因冲突/重复被更新,旧内容会以轻量版本对象 ArchivedTextualMemory(item.py:49)存进原节点 metadata 的 history 列表(item.py:125),同时可用一个 archived_memory_id 指向一个保存完整旧信息(含 sources/embedding)的 archived 节点。

ArchivedTextualMemory 的关键字段:

字段作用定义
version版本号,与活跃记忆的 version 比对item.py:61
update_type为何归档:conflict/duplicate/extract/unrelated/feedbackitem.py:72
archived_memory_id指向存完整旧信息的 archived 节点item.py:76
memory该历史版本的正文快照item.py:69

活跃元数据侧对应有 versionhistoryevolve_to 三个字段(item.py:117-128),共同支撑"版本可追溯 + 可回滚"。

5.3 drop:备份后清库,并滚动保留

TreeTextMemory.drop(tree.py:478)是"清空整库"的安全版:

# 示意,非源码 —— drop 的流程(tree.py:478)
backup_dir = 临时目录 / f"memos_backup_{时间戳}" # 带时间戳的版本化备份
self.dump(backup_dir) # 先把整图导出成 JSON
self._cleanup_old_backups(backup_root, keep_last_n) # 只留最近 keep_last_n 份
self.graph_store.drop_database() # 再真正 drop 库

_cleanup_old_backups(tree.py:505)按目录名里的时间戳倒序排,backups[keep_last_n:] 之外的旧备份 shutil.rmtree 删掉——保证备份目录不会无限膨胀(默认留 30 份)。


6. 后台层级重构:把碎片聚成一棵树

前面(第 4 节)是"来一个整一个"的增量整理。本节是周期性的全局重构:GraphStructureReorganizer 起了两条后台线程(reorganizer.py:96-106),第二条每 100 秒触发一次 optimize_structure(reorganizer.py:158),把平铺的节点聚类、生成"摘要父节点",给图长出层级。

只有 is_reorganize=True(配置 reorganize)时这两条线程才启动。

6.1 全景:optimize_structure 三步走

optimize_structure(scope) reorganizer.py:211


① 取候选:get_structure_optimization_candidates 只挑"孤立/无父无子"的 activated 节点
│ neo4j.py:1544
▼ (不足 min_group_size=20 就跳过)
② 分区:_partition 嵌入向量 KMeans 递归聚类
│ reorganizer.py:459

③ 每个簇:_process_cluster_and_write reorganizer.py:314
├─ 大簇再 _local_subcluster(LLM 切分)
├─ _summarize_cluster → 生成摘要父节点 reorganizer.py:550
├─ _create_parent_node + _link_cluster_nodes(连 PARENT 边)
└─ relation_detector.process_node(见 4.3,当前空转)

6.2 分区:嵌入向量的递归 KMeans

_partition(reorganizer.py:459)的策略很直接:

  • 节点数 ≤ max_cluster_size(20)→ 整堆算一个簇,不聚类。
  • 否则用 MiniBatchKMeans(reorganizer.py:508)按 embedding 聚成 k = ceil(N/20) 个簇,递归地对仍然过大的子簇继续切,直到每簇 ≤ 20。
  • 没有 embedding 的节点被塞进最大的那个簇兜底(reorganizer.py:516)。
  • 最后只保留 size > min_cluster_size(10)的簇。

6.3 摘要父节点:让 LLM 给一簇起个名

_summarize_cluster(reorganizer.py:550)是"长出层级"的关键——它把一簇节点的 key/value/background 拼给 LLM(REORGANIZE_PROMPT),让它产出一个概括性的父节点:

# 示意,非源码 —— _summarize_cluster 的产出(reorganizer.py:571)
parent = GraphDBNode(
memory=parent_value, # LLM 概括出的父记忆正文
metadata=TreeNodeTextualMemoryMetadata(
memory_type=scope,
key=parent_key, # 父节点标题
background=parent_background, # 摘要
sources=build_summary_parent_node(cluster_nodes), # 记下它由哪些子节点汇成
confidence=0.66,
type="topic", # 标成"主题"节点
),
)

生成父节点后:_create_parent_node 写入图(reorganizer.py:609),_link_cluster_nodes 给"父 → 每个子"连 PARENT 边(reorganizer.py:620)。若一簇很大先被 _local_subcluster 切成多个子簇,则形成两层:大簇父 → 各子簇父 → 叶子(reorganizer.py:337-343)。这样原本平铺的记忆就有了 topic → concept → fact 的层级,供检索侧(第 5 章)按层重排。

6.4 消息队列的调度细节

两条后台线程与队列的配合:

  • 优先级队列:QueueMessage.__lt__(reorganizer.py:67)定义 add/remove 优先于 mergemerge 优先于 end——保证 end(停止信号)总排最后被处理。
  • 消费循环 _run_message_consumer_loop(reorganizer.py:133):取消息 → _preprocess_message 把 id 换成真实节点、丢掉已不存在的节点(reorganizer.py:635)→ handle_add 做去重消解。
  • 优雅停止 stop(reorganizer.py:173):塞一条 op="end" 让消费线程退出,再置停调度线程,MemoryManager.close/__del__(manager.py:556-561)会调用它。
  • 等待收尾 wait_until_current_task_done(reorganizer.py:111):queue.join() + 轮询 _is_optimizing 标志,给"写完要立刻检索"的场景一个同步点。

7. 图数据库后端抽象:同一套接口,四种落地

上面所有 graph_store.xxx 调用,都打在一个抽象接口上,底下可换后端。

7.1 接口与工厂

BaseGraphDB(graph_dbs/base.py:11)是抽象基类,规定了节点/边/查询/搜索/结构维护一整套抽象方法(add_nodeadd_edgesearch_by_embeddingget_structure_optimization_candidatesadd_nodes_batch 等)。GraphStoreFactory(graph_dbs/factory.py:11)按 backend 字段选实现:

backend实现类文件
neo4jNeo4jGraphDBgraph_dbs/neo4j.py:99
neo4j-communityNeo4jCommunityGraphDB(继承 Neo4j 版)graph_dbs/neo4j_community.py:22
postgresPostgresGraphDBgraph_dbs/postgres.py:43
polardbPolarDBGraphDBgraph_dbs/polardb.py:101

节点对象统一用 GraphDBNode(graph_dbs/item.py:10)——它直接继承 TextualMemoryItem,所以"记忆项"和"图节点"是同一套字段;边用 GraphDBEdge(graph_dbs/item.py:14)。

7.2 Neo4j 实现里值得记住的几处

以参考实现 Neo4jGraphDB 为例:

  • 两种租户模式(neo4j.py:110 的 docstring):use_multi_db=True 每租户一个独立数据库;use_multi_db=False 所有租户共库,靠每条查询强制 user_name 过滤在节点级隔离。
  • 写入即 upsert:add_node(neo4j.py:224)用 MERGE (n:Memory {id}) + SET n += $metadata,并把 created_at/updated_at 存成 Neo4j datetime;写前经 _prepare_node_metadata/_sanitize_neo4j_metadata 把嵌套结构拍平以适配 Neo4j 的扁平属性模型。
  • 批量写:add_nodes_batch(neo4j.py:273)一次 UNWIND 多个节点。
  • 向量预过滤搜索:search_by_embedding(neo4j.py:840)在 Neo4j ≥ 5.18 上先 WHERE 收窄候选(scope/status/user_name),再用 vector.similarity.cosine() 算分,避免全局 top-k 把过滤条件外的节点挤掉。
  • FIFO 淘汰:remove_oldest_memory(neo4j.py:197)按 updated_at DESC SKIP keep_latest DETACH DELETE,工作记忆的容量上限就靠它维持。
  • 重构候选:get_structure_optimization_candidates(neo4j.py:1544)只选"既无父边也无子边"的 activated 节点,正好喂给第 6 节的聚类。

一个架构判断: 抽象接口里挂着 deduplicate_nodes/detect_conflicts/merge_nodes 三个方法,但 Neo4j 实现里它们都是 raise NotImplementedError(neo4j.py:1201-1225)。去重/冲突消解并不在数据库层做,而是上移到了应用层的 NodeHandler(第 4 节)——因为判定"重复/矛盾"需要 LLM 语义判断,不是 Cypher 能表达的。这是 MemOS 有意的分工。


8. 边界与局限(诚实清单)

  • 两条写入路径行为不一致:批量路径不双写工作记忆(manager.py:239 停用),并行路径才双写;想要工作记忆同步的场景需注意走哪条。
  • 批量 add 的去重只覆盖首节点:handle_add 只取 after_node[0](reorganizer.py:194),批量消息里其余节点这一轮不做增量去重。
  • 丰富关系推理线休眠:RelationAndReasoningDetector.process_node 的因果/推断/时序/聚合四段被注释停用(relation_reason_detector.py:49-80),当前只有 NodeHandler 的 redundant/contradictory 在生效。
  • 去重/冲突全靠 LLM:检测和消解都要调 LLM(handler.py),质量与成本受模型影响;LLM 判 independent 就不会合并,阈值 0.8 以下的相似项也不进候选。
  • 重构是"尽力而为":optimize_structure 有 600s 看门狗(reorganizer.py:229),超时会中断并取消未完成的簇;节点不足 20 直接跳过。
  • 结构脚手架未接线:_ensure_structure_path/_inherit_edges 已实现但未被 add 主链路调用。

9. 横向对比

同 shelf 的其它记忆/上下文系统(参见总库 doc 的对应原理),对"去重与冲突"取舍各异:

  • 本章(MemOS):异步、图上、LLM 语义判定 + 归档保版本;写入快、整理慢养,适合长期演化的知识库。
  • 相邻章节里,03 负责"从原始输入抽取成项",本章负责"项落图并养护",05 负责"从图里召回",06 负责"异步摄取与激活记忆刷新的调度"。四章合起来是 TreeTextMemory 的完整生命周期。

10. 代码地图(导航索引)

主题文件路径关键符号
明文记忆门面 / 对外方法src/memos/memories/textual/tree.pyTreeTextMemory.addreplace_working_memoryget_working_memorysoft_deletedrop_cleanup_old_backups
写入编排src/memos/memories/textual/tree_text_memory/organize/manager.pyMemoryManager.add_add_memories_parallel_add_memories_batch_process_memory_add_to_graph_memory_add_memory_to_db_ensure_structure_path_inherit_edgesextract_working_binding_ids
后台重构守护src/memos/memories/textual/tree_text_memory/organize/reorganizer.pyGraphStructureReorganizerQueueMessage_run_message_consumer_loop_run_structure_organizer_loophandle_addoptimize_structure_partition_summarize_cluster_process_cluster_and_write_link_cluster_nodes
去重 / 冲突消解src/memos/memories/textual/tree_text_memory/organize/handler.pyNodeHandler.detectresolve_hard_update_resolve_in_graph_merge_metadataEMBEDDING_THRESHOLD
关系 / 推理探测(休眠)src/memos/memories/textual/tree_text_memory/organize/relation_reason_detector.pyRelationAndReasoningDetector.process_node_detect_pairwise_causal_condition_relations_infer_fact_nodes_from_relations_parse_relation_result
记忆项 / 版本 / 元数据src/memos/memories/textual/item.pyTextualMemoryItemTreeNodeTextualMemoryMetadataArchivedTextualMemorySourceMessage
图 DB 抽象 + 工厂src/memos/graph_dbs/base.pyfactory.pyitem.pyBaseGraphDBGraphStoreFactoryGraphDBNodeGraphDBEdge
图 DB 后端实现src/memos/graph_dbs/neo4j.py(及 neo4j_community.py/postgres.py/polardb.py)Neo4jGraphDB.add_nodeadd_nodes_batchsearch_by_embeddingremove_oldest_memoryget_structure_optimization_candidatesupdate_nodedrop_database