在 AI Agent 系统中,长期记忆(Long-term Memory)是实现持续性对话、个性化服务、自主决策的核心基础设施。当前大多数 Agent 框架将"长期记忆"简单等同于向量数据库检索(RAG),这种理解在生产环境中会导致严重的记忆碎片化、时序因果混乱和知识更新滞后问题。

本文将从工程实现层面,深入分析 AI Agent 长期记忆系统的分层架构、存储引擎选择、检索与写入策略,以及如何构建具备时序一致性的认知记忆系统。通过完整的代码示例,展示如何在生产环境中实现一个支持 CRUD 操作(创建、检索、更新、删除)的长期记忆模块。


一、记忆系统的层次化模型

人类认知科学将记忆分为感觉记忆、短时记忆和长时记忆(语义记忆与情景记忆)。借鉴这一分层结构,AI Agent 的记忆系统应该设计为四层架构:

┌─────────────────────────────────────────────────────────┐
│  L1: 上下文窗口(Context Window)                          │
│     - 当前会话的原始对话历史,受 token 限制                    │
│     - 工作机制:滑动窗口 + 重要度评分截断                       │
└─────────────────────────────────────────────────────────┘
                            ↕ 压缩与摘要
┌─────────────────────────────────────────────────────────┐
│  L2: 工作记忆(Working Memory)                            │
│     - 当前任务的推理上下文(目标、计划、中间结论)               │
│     - 生命周期:单次任务会话内                                │
│     - 数据结构:结构化 JSON/Tree 表示                        │
└─────────────────────────────────────────────────────────┘
                            ↕ 重要事件写入 / 触发式召回
┌─────────────────────────────────────────────────────────┐
│  L3: 情景记忆(Episodic Memory)                           │
│     - 带时间戳的交互历史片段                                  │
│     - 特性:时序排序、因果关联、衰减评分                       │
│     - 检索方式:时间范围查询 + 语义搜索混合                    │
└─────────────────────────────────────────────────────────┘
                            ↕ 模式提取与知识蒸馏
┌─────────────────────────────────────────────────────────┐
│  L4: 语义记忆(Semantic Memory)                           │
│     - 结构化知识图谱 / 事实集合                              │
│     - 特性:去重、版本化、置信度评分                          │
│     - 检索方式:语义相似度 + 精确匹配混合                     │
└─────────────────────────────────────────────────────────┘

这种四层分离的核心价值在于:每层使用不同的存储引擎和检索策略,避免用单一技术解决所有问题。L1 和 L2 通常驻留在内存中,而 L3 和 L4 必须持久化到磁盘或外部存储服务。


二、存储引擎选型与数据模型

2.1 情景记忆的存储方案

情景记忆的核心需求是按时间范围查询、按语义检索、支持 TTL 衰减。推荐使用 PostgreSQL + pgvector 的组合:

-- 情景记忆表结构
CREATE TABLE episodic_memories (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    agent_id TEXT NOT NULL,
    session_id TEXT NOT NULL,
    content TEXT NOT NULL,
    embedding vector(1536),                    -- OpenAI text-embedding-3-small
    importance FLOAT DEFAULT 0.5,              -- 重要度评分 0-1
    created_at TIMESTAMPTZ DEFAULT NOW(),
    accessed_count INT DEFAULT 0,
    last_accessed_at TIMESTAMPTZ,
    metadata JSONB DEFAULT '{}',               -- 来源、情绪、实体等
    ttl TIMESTAMPTZ,                           -- 可选过期时间
    INDEX idx_episodic_time (agent_id, created_at DESC),
    INDEX idx_episodic_importance (agent_id, importance DESC)
);

-- HNSW 索引用于向量检索
CREATE INDEX ON episodic_memories
    USING hnsw (embedding vector_cosine_ops)
    WITH (m = 16, ef_construction = 200);

选择 PostgreSQL 而非专用向量数据库的原因:情景记忆需要强大的事务支持——当 Agent 在一个涉及多步骤的任务中需要原子性写入多个记忆片段,以及时间序列查询的复杂条件过滤时,PostgreSQL 的成熟生态是最佳选择。

2.2 语义记忆的图模型

语义记忆更适合图数据库来存储实体间的关系网络,我们可以使用 Apache AGE(基于 PostgreSQL 的图扩展)来实现:

-- 知识节点表
CREATE TABLE semantic_knowledge (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    entity_type TEXT NOT NULL,     -- Person / Concept / Fact / Skill
    entity_name TEXT NOT NULL,
    properties JSONB DEFAULT '{}',
    confidence FLOAT DEFAULT 0.8,
    version INT DEFAULT 1,
    source TEXT,                   -- 来源记忆 ID
    created_at TIMESTAMPTZ DEFAULT NOW(),
    updated_at TIMESTAMPTZ DEFAULT NOW(),
    UNIQUE(entity_type, entity_name)
);

-- 知识关系表
CREATE TABLE knowledge_relations (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    from_entity_id UUID REFERENCES semantic_knowledge(id),
    to_entity_id UUID REFERENCES semantic_knowledge(id),
    relation_type TEXT NOT NULL,   -- knows / depends_on / contradicts / implies
    weight FLOAT DEFAULT 1.0,
    properties JSONB DEFAULT '{}',
    created_at TIMESTAMPTZ DEFAULT NOW()
);

CREATE INDEX idx_relations_from ON knowledge_relations(from_entity_id, relation_type);
CREATE INDEX idx_relations_to ON knowledge_relations(to_entity_id, relation_type);

2.3 Rust 实现的记忆存储管理器

下面是一个用 Rust 实现的核心记忆管理模块,封装了四层存储引擎的读写操作:

use sqlx::{postgres::PgPoolOptions, Pool, Postgres};
use serde::{Deserialize, Serialize};
use chrono::{DateTime, Utc, Duration};
use uuid::Uuid;

/// 记忆层级分类
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum MemoryTier {
    Working,    // L2: 工作记忆(内存中)
    Episodic,   // L3: 情景记忆(PG 持久化)
    Semantic,   // L4: 语义记忆(图结构)
}

/// 记忆条目
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MemoryEntry {
    pub id: Option<String>,
    pub tier: MemoryTier,
    pub agent_id: String,
    pub content: String,
    pub embedding: Option<Vec<f32>>,
    pub importance: f32,
    pub metadata: serde_json::Value,
    pub created_at: Option<DateTime<Utc>>,
}

/// 检索结果,包含相似度得分
#[derive(Debug, Clone)]
pub struct MemoryResult {
    pub entry: MemoryEntry,
    pub similarity: f32,
    pub recency_boost: f32,    // 时间衰减补偿
}

/// 长期记忆管理器
pub struct LongTermMemoryManager {
    db_pool: Pool<Postgres>,
    embedding_client: reqwest::Client,
    embedding_endpoint: String,
    api_key: String,
}

impl LongTermMemoryManager {
    pub async fn new(database_url: &str, api_key: &str) -> Result<Self, Box<dyn std::error::Error>> {
        let pool = PgPoolOptions::new()
            .max_connections(10)
            .connect(database_url)
            .await?;

        Ok(Self {
            db_pool: pool,
            embedding_client: reqwest::Client::new(),
            embedding_endpoint: "https://api.openai.com/v1/embeddings".to_string(),
            api_key: api_key.to_string(),
        })
    }

    /// 生成文本向量
    async fn embed(&self, text: &str) -> Result<Vec<f32>, Box<dyn std::error::Error>> {
        let response = self.embedding_client
            .post(&self.embedding_endpoint)
            .header("Authorization", format!("Bearer {}", self.api_key))
            .json(&serde_json::json!({
                "model": "text-embedding-3-small",
                "input": text
            }))
            .send().await?
            .json::<serde_json::Value>().await?;

        let embedding = response["data"][0]["embedding"]
            .as_array()
            .unwrap()
            .iter()
            .map(|v| v.as_f64().unwrap() as f32)
            .collect();

        Ok(embedding)
    }

    /// 写入情景记忆
    pub async fn store_episodic(
        &self,
        agent_id: &str,
        session_id: &str,
        content: &str,
        importance: f32,
    ) -> Result<String, Box<dyn std::error::Error>> {
        let embedding = self.embed(content).await?;
        let embedding_str = format!("[{}]", embedding.iter().map(|f| f.to_string()).collect::<Vec<_>>().join(","));
        let id = Uuid::new_v4().to_string();

        sqlx::query(
            r#"
            INSERT INTO episodic_memories (id, agent_id, session_id, content, embedding, importance)
            VALUES ($1, $2, $3, $4, $5::vector, $6)
            "#
        )
        .bind(&id)
        .bind(agent_id)
        .bind(session_id)
        .bind(content)
        .bind(&embedding_str)
        .bind(importance)
        .execute(&self.db_pool)
        .await?;

        Ok(id)
    }

    /// 混合检索:语义相似度 + 时间衰减 + 重要度加权
    pub async fn recall(
        &self,
        agent_id: &str,
        query: &str,
        limit: usize,
        time_window_days: Option<i64>,
    ) -> Result<Vec<MemoryResult>, Box<dyn std::error::Error>> {
        let query_embedding = self.embed(query).await?;
        let emb_str = format!("[{}]", query_embedding.iter().map(|f| f.to_string()).collect::<Vec<_>>().join(","));

        let time_filter = time_window_days.map(|days| {
            format!("AND created_at > NOW() - INTERVAL '{} days'", days)
        }).unwrap_or_default();

        let rows = sqlx::query_as::<_, EpisodicRow>(
            &format!(r#"
                SELECT id, agent_id, session_id, content, importance, created_at, metadata,
                       (embedding <=> $1::vector) as distance
                FROM episodic_memories
                WHERE agent_id = $2 {}
                ORDER BY (
                    0.6 * (1.0 - (embedding <=> $1::vector)) +  -- 语义相似度
                    0.25 * importance +                             -- 重要度
                    0.15 * EXTRACT(EPOCH FROM (NOW() - created_at)) / 86400.0 * -1.0  -- 时间衰减
                ) DESC
                LIMIT $3
            "#, time_filter)
        )
        .bind(&emb_str)
        .bind(agent_id)
        .bind(limit as i64)
        .fetch_all(&self.db_pool)
        .await?;

        let results: Vec<MemoryResult> = rows.into_iter().map(|r| {
            let days_since = (Utc::now() - r.created_at.unwrap_or_else(Utc::now)).num_days() as f32;
            let recency_boost = 1.0 / (1.0 + days_since * 0.1);  // 时间衰减因子

            MemoryResult {
                entry: MemoryEntry {
                    id: Some(r.id),
                    tier: MemoryTier::Episodic,
                    agent_id: r.agent_id,
                    content: r.content,
                    embedding: None,
                    importance: r.importance,
                    metadata: r.metadata.unwrap_or(serde_json::json!({})),
                    created_at: r.created_at,
                },
                similarity: 1.0 - r.distance,
                recency_boost,
            }
        }).collect();

        Ok(results)
    }

    /// 记忆巩固:将碎片化情景记忆提炼为结构化语义知识
    pub async fn consolidate(&self, agent_id: &str) -> Result<usize, Box<dyn std::error::Error>> {
        // 获取最近的高重要度记忆片段
        let memories = sqlx::query_as::<_, EpisodicRow>(
            r#"
            SELECT id, content, importance
            FROM episodic_memories
            WHERE agent_id = $1
              AND importance > 0.7
              AND created_at > NOW() - INTERVAL '7 days'
            ORDER BY importance DESC
            LIMIT 50
            "#
        )
        .bind(agent_id)
        .fetch_all(&self.db_pool)
        .await?;

        // 此处调用 LLM 进行知识提取(实体识别、关系抽取)
        // 简化为示例:将记忆内容直接写入语义知识表
        let mut consolidated = 0;
        for mem in &memories {
            sqlx::query(
                r#"
                INSERT INTO semantic_knowledge (entity_type, entity_name, properties, confidence, source)
                VALUES ('Fact', $1, $2, $3, $4)
                ON CONFLICT (entity_type, entity_name)
                DO UPDATE SET
                    confidence = GREATEST(semantic_knowledge.confidence, EXCLUDED.confidence),
                    version = semantic_knowledge.version + 1,
                    updated_at = NOW()
                "#
            )
            .bind(&mem.content.chars().take(200).collect::<String>())
            .bind(&serde_json::json!({"original_id": mem.id}))
            .bind(mem.importance)
            .bind(&mem.id)
            .execute(&self.db_pool)
            .await?;
            consolidated += 1;
        }

        Ok(consolidated as usize)
    }
}

/// 数据库行映射(简化版)
#[derive(sqlx::FromRow, Debug)]
struct EpisodicRow {
    id: String,
    agent_id: String,
    session_id: Option<String>,
    content: String,
    importance: f32,
    created_at: Option<DateTime<Utc>>,
    metadata: Option<serde_json::Value>,
    distance: f32,
}

三、检索排序的工程细节

单纯的向量余弦相似度往往产生语义正确但上下文不相关的检索结果。生产环境中的检索排序需要综合考虑以下几个因子:

3.1 时间衰减模型

记忆的价值随时间递减,但重要事件(如一个反复出现的系统故障)应衰减更慢。推荐使用指数衰减 + 重要度调节因子:

import math
from datetime import datetime, timezone

def compute_memory_score(
    cosine_similarity: float,
    importance: float,
    created_at: datetime,
    accessed_count: int,
    query_context: dict
) -> float:
    """
    综合评分函数:结合语义相关度、重要度、时间衰减和访问频率

    权重分配:
    - 语义相似度: 0.45
    - 重要度: 0.25
    - 时间衰减: 0.20
    - 访问频率: 0.10
    """
    # 基础语义分(cosine 距离转相似度)
    semantic_score = max(0.0, cosine_similarity)

    # 重要度有最小值保护,高重要记忆永远不会被完全淹没
    importance_score = min(1.0, importance)

    # 时间衰减:半衰期与重要度关联(重要记忆衰减更慢)
    now = datetime.now(timezone.utc)
    age_days = max(0, (now - created_at).total_seconds() / 86400)
    half_life = 7 + importance_score * 21  # 重要度 0.0 → 7天半衰期;1.0 → 28天
    time_score = math.exp(-0.693 * age_days / half_life)  # ln(2) ≈ 0.693

    # 访问频率:被反复访问的记忆更可能是"核心知识"
    frequency_score = min(1.0, math.log(1 + accessed_count) / math.log(20))

    # 加权融合
    final_score = (
        0.45 * semantic_score +
        0.25 * importance_score +
        0.20 * time_score +
        0.10 * frequency_score
    )

    return final_score

3.2 重排序与去重

向量检索返回 Top-K 候选后,还需要执行去重和重排序:

def rerank_and_deduplicate(candidates: list[MemoryResult], threshold: float = 0.85) -> list[MemoryResult]:
    """
    1. 按综合评分排序
    2. 对高度相似的结果去重(保留评分最高者)
    3. 返回最终列表
    """
    # 按评分降序排列
    candidates.sort(key=lambda r: r.final_score, reverse=True)

    selected = []
    for candidate in candidates:
        # 检查是否与已选结果重复
        is_duplicate = False
        for kept in selected:
            # 使用 Jaccard 相似度快速判断内容重叠
            overlap = _jaccard_similarity(candidate.entry.content, kept.entry.content)
            if overlap > threshold:
                is_duplicate = True
                break

        if not is_duplicate:
            selected.append(candidate)

    return selected

def _jaccard_similarity(text_a: str, text_b: str) -> float:
    """基于 N-gram token 的 Jaccard 相似度"""
    def ngrams(text: str, n: int = 3) -> set:
        tokens = text.lower().split()
        return set(tuple(tokens[i:i+n]) for i in range(len(tokens)-n+1))

    set_a = ngrams(text_a)
    set_b = ngrams(text_b)
    if not set_a or not set_b:
        return 0.0
    return len(set_a & set_b) / len(set_a | set_b)

四、记忆更新与遗忘的一致性保障

长期记忆最大的工程挑战不是写入和读取,而是更新一致性和遗忘策略。

4.1 事件驱动的记忆更新

当 Agent 发现之前的记忆有误(比如用户纠正了一个偏好),系统需要能够定位并更新相关记忆。核心思路是将 LLM 的"意图识别"与 CRUD 操作结合:

from enum import Enum

class MemoryOp(Enum):
    CREATE = "create"
    UPDATE = "update"
    DELETE = "delete"
    CONSOLIDATE = "consolidate"  # 合并碎片记忆

class MemoryUpdateEngine:
    def __init__(self, storage_manager, llm_client):
        self.storage = storage_manager
        self.llm = llm_client

    async def process_memory_intent(
        self, 
        agent_id: str, 
        user_message: str, 
        context_messages: list[str]
    ) -> MemoryOp:
        """
        分析用户消息,判断是否涉及记忆更新操作
        """
        system_prompt = """分析用户消息,判断是否需要操作记忆系统。
        可能的意图:
        - UPDATE: 用户说"其实我喜欢...""不对,之前说错了...",表明需要纠正
        - DELETE: 用户说"忘掉...""别记住...", 表明需要删除
        - CREATE: 新事实或偏好,不需要覆盖已有内容
        - NONE: 普通对话,不涉及记忆

        只返回操作类型和相关的记忆ID(如果有)"""

        response = await self.llm.chat_json(
            system=system_prompt,
            messages=context_messages + [user_message],
            schema=MemoryDecision
        )
        return response

    async def update_memory(
        self,
        agent_id: str,
        target_memory_id: str,
        correction: str,
    ) -> dict:
        """
        原子操作:更新记忆并记录变更历史
        """
        # 1. 读取原始记忆
        original = await self.storage.get_by_id(target_memory_id)

        # 2. 写入新版本
        new_id = await self.storage.store_episodic(
            agent_id=agent_id,
            session_id="system_update",
            content=correction,
            importance=original.importance,
        )

        # 3. 标记旧版本为已取代
        await self.storage.mark_superseded(target_memory_id, new_id)

        return {"old_id": target_memory_id, "new_id": new_id}

4.2 遗忘曲线与存储成本控制

无限制的存储增长是不可行的。基于艾宾浩斯遗忘曲线的变体,我们可以设计一个渐进式的记忆衰减策略:

  • 立即访问后: 记忆强度设为 1.0
  • 第 N 次复习: 衰减系数按照 interval = previous_interval × EF(EF 为易度因子,范围 1.3-2.5)
  • 未复习超过阈值: 进入"压缩区"——只保存摘要而非原文
  • 超过最大周期: 归档到冷存储(S3 等),不再参与实时召回
class ForgettingCurveEviction:
    """基于改造版 SM-2 算法的记忆淘汰"""

    def __init__(self, 
                 max_memories: int = 50_000,
                 compress_after_days: int = 30,
                 archive_after_days: int = 90):
        self.max_memories = max_memories
        self.compress_after = compress_after_days
        self.archive_after = archive_after_days

    async def run_eviction(self, pool, agent_id: str):
        """执行一轮淘汰,返回被处理的记忆数量"""
        total_count = await self._count_memories(pool, agent_id)

        if total_count <= self.max_memories:
            return 0

        excess = total_count - self.max_memories
        archived = 0
        compressed = 0
        deleted = 0

        # 第一步:归档超期记忆(转移到 S3 冷存储)
        archive_threshold = datetime.now(timezone.utc) - timedelta(days=self.archive_after)
        rows = await pool.fetch(
            """UPDATE episodic_memories 
               SET archived_at = NOW(), embedding = NULL
               WHERE id IN (
                   SELECT id FROM episodic_memories
                   WHERE agent_id = $1 
                     AND created_at < $2
                     AND archived_at IS NULL
                   ORDER BY importance ASC, last_accessed_at ASC NULLS FIRST
                   LIMIT $3
               ) RETURNING id""",
            agent_id, archive_threshold, excess
        )
        archived = len(rows)

        # 第二步:压缩近期但不常访问的记忆
        if archived < excess:
            remaining = excess - archived
            compress_threshold = datetime.now(timezone.utc) - timedelta(days=self.compress_after)
            # 调用 LLM 生成摘要替换原文
            rows = await pool.fetch(
                """SELECT id, content FROM episodic_memories
                   WHERE agent_id = $1
                     AND created_at < $2
                     AND accessed_count < 3
                     AND archived_at IS NULL
                     AND compressed = FALSE
                   ORDER BY importance ASC
                   LIMIT $3""",
                agent_id, compress_threshold, remaining
            )
            for row in rows:
                summary = await self._compress_content(row['content'])
                await pool.execute(
                    "UPDATE episodic_memories SET content = $1, compressed = TRUE WHERE id = $2",
                    summary, row['id']
                )
                compressed += 1

        return {"archived": archived, "compressed": compressed, "total_processed": archived + compressed}

五、生产部署的架构建议

将上述组件组装为一个完整的生产记忆系统,典型的架构如下:

                    ┌──────────────────────────┐
                    │      AI Agent Runtime      │
                    │  (LangChain / Custom)      │
                    └─────────┬────────────────┘
                              │ Recall / Write / Update
                              ▼
                    ┌──────────────────────────┐
                    │    Memory Gateway API     │
                    │  (Go/gRPC, <5ms P99)      │
                    │  - 请求路由 & 限流          │
                    │  - A/B 实验分流             │
                    └─────────┬────────────────┘
                              │
              ┌───────────────┼───────────────────┐
              ▼               ▼                   ▼
    ┌─────────────────┐ ┌──────────┐ ┌──────────────────────┐
    │ PostgreSQL +    │ │ Redis     │ │ LLM Embedding        │
    │ pgvector        │ │ Cluster  │ │ Service              │
    │ (主存储,        │ │ (L1/L2   │ │ (batch + streaming   │
    │  强一致性)       │ │  热缓存)  │ │  双通道)             │
    └─────────────────┘ └──────────┘ └──────────────────────┘
              │
              ▼ (定期归档)
    ┌─────────────────┐
    │ S3 / HDFS       │
    │ (冷存储, 摘要)   │
    └─────────────────┘

关键技术指标:
- API P99 延迟: < 50ms(不含 LLM 调用)
- 写入吞吐: > 500 QPS(单 PG 节点)
- 召回相关性: NDCG@10 > 0.85(通过人工评测集)
- 存储成本: 每百万记忆 < $3/月(压缩+归档后)

运维建议:
1. 对 embedding 向量列执行 VACUUM 定期维护,避免 HNSW 索引性能退化
2. 监控"记忆效用率"——即被实际召回并用于推理的记忆占比,目标 > 60%
3. 建立记忆质量反馈回路:LLM 在推理时使用某条记忆后,异步更新其重要性评分


六、总结

AI Agent 的长期记忆系统不是一个"买了向量数据库就能用"的简单中间件。它需要像对待数据库一样认真设计数据模型,像对待缓存一样管理分层存储,像对待知识图谱一样维护关系一致性。

工程化的核心原则是:

  1. 分层而非混合:上下文窗口、工作记忆、情景记忆、语义记忆各司其职
  2. 检索不是全部:更新一致性、遗忘策略、知识蒸馏同样关键
  3. 没有免费午餐的嵌入:选择 embedding 模型时需要在速度、成本、维度之间权衡
  4. 闭环监控:记忆系统的价值最终通过 Agent 任务成功率来衡量,而非检索准确率

通过本文的分层架构和代码示例,开发者可以在一周内构建一个生产可用的 Agent 长期记忆模块,为打造具备"持续性记忆"能力的 AI Agent 打下坚实基础。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部