Agent 记忆系统:三层架构的完整实现
一、为什么不能只靠 Checkpoint
LangGraph 自带的 Checkpoint 能记住对话历史,但有两个问题:
- 只在进程内有效。服务重启,内存里的 checkpoint 全丢(除非配了持久化后端)。
- 对话长了会爆。Supervisor 每次调 LLM 都要带完整
messages列表,20 轮就是 15000+ token。
所以我搞了三层记忆,各管各的时间尺度:
L1 瞬时记忆 → LangGraph Checkpoint → 单次图执行内的状态流转
L2 短期记忆 → SQLite → 跨请求、跨重启的对话历史
L3 长期记忆 → SQLite 用户画像 → 用户的偏好、历史目的地
三层之间不是互相替代,是各司其职。下面一层层拆开讲。
二、L1 瞬时记忆:AgentState 的设计
from typing import TypedDict, Annotated
from langgraph.graph.message import add_messages
class AgentState(TypedDict):
# ---- 消息列表(带 Reducer) ----
messages: Annotated[list, add_messages] # 自动追加 + 去重
# ---- 用户输入 ----
user_input: str
# ---- Supervisor 提取的结构化参数 ----
destination: str
days: int
budget: str
departure: str
transport_mode: str
preferences: str
# ---- 子 Agent 的产出(并行写入,互不冲突) ----
route_result: str
hotel_result: str
food_result: str
info_result: str
# ---- 汇总结果 ----
final_result: str
# ---- 控制字段 ----
agents_to_call: list[str] # Supervisor 决定调用哪些 Agent
is_casual: bool # 闲聊标记,跳过规划流程
need_clarify: bool # 需要追问用户
clarify_msg: str # 追问内容
# ---- 摘要状态 ----
running_summary: dict # langmem RunningSummary 序列化
# ---- 会话标识 ----
user_id: str
session_id: str
request_id: str
关键设计就一条:每个子 Agent 只写自己的 _result 字段。 并行执行时 route 写 route_result、hotel 写 hotel_result——各写各的,不会冲突。LangGraph 在 fan-in 时自动合并,这是它能稳定并行跑的底层保证。
messages 字段用了 add_messages Reducer,它的合并规则:
# 伪代码:add_messages 的实际行为
def add_messages(old: list, new: list) -> list:
merged = {m.id: m for m in old}
for m in new:
if m.id in merged:
merged[m.id] = m # 同 ID → 更新(覆盖)
else:
merged[m.id] = m # 新 ID → 追加
return list(merged.values())
三、L2 短期记忆:SQLite 对话存储
Checkpoint 的问题是不持久、不跨进程。所以 Agent 每次说完话,额外写一份到 SQLite。
3.1 表结构
CREATE TABLE conversations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
user_id TEXT NOT NULL DEFAULT 'guest',
session_id TEXT NOT NULL,
role TEXT NOT NULL, -- 'user' | 'assistant' | 'agent'
content TEXT NOT NULL,
agent TEXT, -- 哪个 Agent 产生的(supervisor/route/hotel/food/info)
metadata TEXT DEFAULT '{}', -- JSON,存额外信息
created_at REAL NOT NULL
);
CREATE TABLE summaries (
id INTEGER PRIMARY KEY AUTOINCREMENT,
user_id TEXT NOT NULL,
session_id TEXT NOT NULL,
summary TEXT NOT NULL,
msg_count INTEGER NOT NULL, -- 这个摘要覆盖了多少条消息
created_at REAL NOT NULL
);
-- 索引:按会话查历史消息
CREATE INDEX idx_conv_session ON conversations(user_id, session_id);
CREATE INDEX idx_conv_created ON conversations(created_at);
3.2 读写操作
import sqlite3
import time
import json
class ConversationStore:
def __init__(self, db_path: str = "data/memory.db"):
self.conn = sqlite3.connect(db_path, check_same_thread=False)
self.conn.execute("PRAGMA journal_mode=WAL") # 读写并发
self._init_tables()
def _init_tables(self):
self.conn.executescript("""
CREATE TABLE IF NOT EXISTS conversations (...);
CREATE TABLE IF NOT EXISTS summaries (...);
CREATE INDEX IF NOT EXISTS idx_conv_session ON conversations(user_id, session_id);
CREATE INDEX IF NOT EXISTS idx_conv_created ON conversations(created_at);
""")
self.conn.commit()
def save_message(self, user_id: str, session_id: str,
role: str, content: str, agent: str = None):
self.conn.execute(
"""INSERT INTO conversations (user_id, session_id, role, content, agent, created_at)
VALUES (?, ?, ?, ?, ?, ?)""",
(user_id, session_id, role, content, agent, time.time())
)
self.conn.commit()
def get_recent_messages(self, user_id: str, session_id: str,
limit: int = 10) -> list[dict]:
rows = self.conn.execute(
"""SELECT role, content, agent, created_at
FROM conversations
WHERE user_id = ? AND session_id = ?
ORDER BY created_at DESC
LIMIT ?""",
(user_id, session_id, limit)
).fetchall()
# 返回时按时间正序(最旧在前,最新在后)
return [
{"role": r[0], "content": r[1], "agent": r[2], "created_at": r[3]}
for r in reversed(rows)
]
def get_message_count(self, user_id: str, session_id: str) -> int:
row = self.conn.execute(
"SELECT COUNT(*) FROM conversations WHERE user_id = ? AND session_id = ?",
(user_id, session_id)
).fetchone()
return row[0] if row else 0
四、上下文组装:每次调 Supervisor 前拼什么
这是记忆系统最核心的函数。每次 Supervisor 执行前,从 SQLite 拉数据、拼成结构化的上下文注入 prompt。
def build_context(user_id: str, session_id: str, store: ConversationStore) -> str:
parts = []
# -------- 1. 用户画像(L3) --------
profile = get_user_profile(user_id) # 见第五节
if profile:
if profile.get("visited_cities"):
recent = profile["visited_cities"][-5:] # 最近 5 个城市
parts.append(f"[用户画像] 常去城市: {', '.join(recent)}")
if profile.get("preferred_budget"):
parts.append(f"[用户画像] 预算偏好: {profile['preferred_budget']}")
if profile.get("preferred_style"):
parts.append(f"[用户画像] 旅行风格: {profile['preferred_style']}")
# -------- 2. 最新摘要 --------
summary = get_latest_summary(user_id, session_id, store)
if summary:
parts.append(f"[对话摘要] {summary}")
# -------- 3. 最近 N 条原始消息 --------
recent = store.get_recent_messages(user_id, session_id, limit=6)
if recent:
lines = []
for m in recent:
role_label = "用户" if m["role"] == "user" else "助手"
# 截断每条到 200 字,避免一条超长消息占满上下文
text = m["content"][:200]
lines.append(f" {role_label}: {text}")
parts.append("[最近对话]\n" + "\n".join(lines))
return "\n".join(parts)
拼出来的上下文大概长这样:
[用户画像] 常去城市: 杭州, 成都, 西安
[用户画像] 预算偏好: 5000元
[用户画像] 旅行风格: 喜欢自然风光
[对话摘要] 用户要求规划成都3天旅行,预算5000元。已推荐宽窄巷子、太古里、锦里等景点...
[最近对话]
用户: 成都三天
助手: 好的,为您规划成都3天行程...
用户: 酒店换便宜点的,最好在武侯区
Supervisor 拿到这份上下文后,能理解"用户之前在聊成都、预算 5000、喜欢自然风光、刚问过换便宜酒店",不需要用户重复说。
五、L3 长期记忆:用户画像
def get_user_profile(user_id: str) -> dict | None:
row = store.conn.execute(
"""SELECT preferred_budget, preferred_style, visited_cities,
last_destination, trip_count
FROM user_profiles WHERE user_id = ?""",
(user_id,)
).fetchone()
if not row:
return None
return {
"preferred_budget": row[0],
"preferred_style": row[1],
"visited_cities": json.loads(row[2] or "[]"),
"last_destination": row[3],
"trip_count": row[4],
}
def update_user_profile(user_id: str, state: dict):
"""每次行程规划完成后更新画像"""
dest = state.get("destination", "")
budget = state.get("budget", "")
profile = get_user_profile(user_id)
if profile:
# 追加新目的地(去重)
cities = profile["visited_cities"]
if dest and dest not in cities:
cities.append(dest)
if len(cities) > 20: # 最多保留 20 个
cities = cities[-20:]
store.conn.execute(
"""UPDATE user_profiles
SET preferred_budget = ?,
visited_cities = ?,
last_destination = ?,
trip_count = trip_count + 1,
updated_at = ?
WHERE user_id = ?""",
(budget or profile["preferred_budget"],
json.dumps(cities, ensure_ascii=False),
dest,
time.time(),
user_id)
)
else:
# 新用户
store.conn.execute(
"""INSERT INTO user_profiles
(user_id, preferred_budget, visited_cities, last_destination, trip_count, updated_at)
VALUES (?, ?, ?, ?, 1, ?)""",
(user_id, budget, json.dumps([dest] if dest else [], ensure_ascii=False),
dest, time.time())
)
store.conn.commit()
更新画像的时机:aggregator 汇总完成后。 这时候 state 里已有完整的 destination、budget 等字段。
六、摘要触发——什么时候该压缩旧对话
SQLite 里的消息越存越多,不可能每次都把完整历史塞进上下文。摘要机制是"消息太多时自动压缩"。
SUMMARY_TRIGGER_COUNT = 8 # 每新增 8 条消息触发一次摘要
KEEP_RECENT = 4 # 始终保留最近 4 条原始消息
def maybe_summarize(user_id: str, session_id: str, llm):
"""检查是否需要触发摘要,如果需要就在后台生成"""
store = ConversationStore()
# 1. 计算自上次摘要以来的新消息数
last_summary = store.conn.execute(
"SELECT msg_count FROM summaries WHERE user_id = ? AND session_id = ? ORDER BY created_at DESC LIMIT 1",
(user_id, session_id)
).fetchone()
last_summarized_count = last_summary[0] if last_summary else 0
current_count = store.get_message_count(user_id, session_id)
new_since_last = current_count - last_summarized_count
# 2. 没到阈值,跳过
if new_since_last < SUMMARY_TRIGGER_COUNT:
return
# 3. 取出需要被摘要的消息(排除最近 KEEP_RECENT 条)
all_msgs = store.get_recent_messages(user_id, session_id, limit=9999)
to_summarize = all_msgs[:-KEEP_RECENT]
if not to_summarize:
return
# 4. 拼成文本,让 LLM 生成摘要
dialog_text = "\n".join(
f"{'用户' if m['role'] == 'user' else '助手'}: {m['content'][:300]}"
for m in to_summarize
)
summary_prompt = f"""将以下对话压缩为不超过 200 字的摘要。
重点保留:
- 目的地、天数、预算
- 用户明确提过的偏好和约束
- Agent 给出的关键推荐
对话:
{dialog_text}
摘要:"""
try:
resp = llm.invoke(summary_prompt, max_tokens=256)
summary_text = resp.content.strip()
# 5. 存摘要
store.conn.execute(
"""INSERT INTO summaries (user_id, session_id, summary, msg_count, created_at)
VALUES (?, ?, ?, ?, ?)""",
(user_id, session_id, summary_text, current_count, time.time())
)
store.conn.commit()
except Exception as e:
# 摘要失败不影响主流程
print(f"摘要生成失败: {e}")
设计要点:
SUMMARY_TRIGGER_COUNT = 8,太小频繁调 LLM,太大上下文膨胀KEEP_RECENT = 4,最近 4 条始终保留原始文本,不压缩——因为最新的对话最可能被追问- 摘要失败不阻塞主流程,只是这轮不压了,下轮再试
七、在 Supervisor 中串起来
def supervisor_node(state: AgentState) -> dict:
user_id = state.get("user_id", "guest")
session_id = state.get("session_id", "default")
store = ConversationStore()
# 1. 保存用户输入(L2)
store.save_message(user_id, session_id, "user", state["user_input"])
# 2. 组装上下文(L2 + L3)
context = build_context(user_id, session_id, store)
# 3. 从 Checkpoint 恢复摘要状态(L1)
running_summary = _deserialize_running_summary(state.get("running_summary"))
# 4. 调 langmem 做消息级摘要(L1 内的压缩)
messages = state.get("messages", [])
try:
result = summarize_messages(
messages=messages,
running_summary=running_summary,
model=supervisor_agent,
max_tokens=4000,
max_tokens_before_summary=3000,
max_summary_tokens=256,
)
processed_messages = result.messages
updated_running_summary = result.running_summary
except Exception:
processed_messages = messages
updated_running_summary = running_summary
# 5. 构建完整 Prompt
system_prompt = SUPERVISOR_SYSTEM_PROMPT
if context:
system_prompt += f"\n\n--- 上下文记忆 ---\n{context}\n---"
# 6. 调 LLM
resp = supervisor_agent.invoke(
[SystemMessage(content=system_prompt)]
+ processed_messages
+ [HumanMessage(content=state["user_input"])]
)
# 7. 保存助手回复(L2)
reply = extract_reply(resp)
store.save_message(user_id, session_id, "assistant", reply, "supervisor")
# 8. 触发后台摘要检查(L2)
maybe_summarize(user_id, session_id, supervisor_agent)
# 9. 更新用户画像(L3)—— 在 aggregator 完成后做,这里只是占位
# update_user_profile(user_id, state) 见 aggregator_node
return {
"messages": [HumanMessage(content=state["user_input"]), AIMessage(content=reply)],
"running_summary": _serialize_running_summary(updated_running_summary),
# ... 其他字段由后续逻辑设置
}
八、三层记忆的协作全景
一次完整请求的数据流:
用户输入 "酒店换便宜点的"
│
├─ save_message() → SQLite conversations 表(L2 写入)
│
├─ build_context():
│ ├─ get_user_profile() → SQLite user_profiles 表(L3 读取)
│ ├─ get_latest_summary() → SQLite summaries 表(L2 读取)
│ └─ get_recent_messages() → SQLite conversations 表(L2 读取)
│
├─ summarize_messages():
│ └─ 从 Checkpoint 恢复 RunningSummary(L1)
│ └─ 如果需要压缩,生成新摘要
│
├─ Supervisor 调用 LLM(带上 L1 + L2 + L3 的上下文)
│
├─ 子 Agent 并行执行 → 各写各的 _result 字段
│
├─ Aggregator 汇总
│ └─ update_user_profile() → SQLite user_profiles 表(L3 写入)
│
├─ save_message() → SQLite conversations 表(L2 写入)
│
└─ maybe_summarize() → 可能触发摘要(L2 压缩)
三层各干各的,互不替代:
- **L1(Checkpoint)**管"这一次图执行内的消息怎么累积"
- **L2(SQLite)**管"跨请求、跨重启的对话历史怎么存和查"
- **L3(画像)**管"这个用户长期来看有什么偏好"
导航
← 上一篇:部署心得与性能优化 | → 下一篇:从单 Agent 到多 Agent:一次真实的架构改造