Agent 记忆系统:三层架构的完整实现

一、为什么不能只靠 Checkpoint

LangGraph 自带的 Checkpoint 能记住对话历史,但有两个问题:

  1. 只在进程内有效。服务重启,内存里的 checkpoint 全丢(除非配了持久化后端)。
  2. 对话长了会爆。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 里已有完整的 destinationbudget 等字段。

六、摘要触发——什么时候该压缩旧对话

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:一次真实的架构改造