Memory Management using Milvus


title: “使用 Milvus 进行内存管理” weight: 500


在本模块中,我们将使用与 RAG with Milvus lab 相同的 Milvus 向量数据库为 agent 提供对话内存。这一次我们存储的是客户自己的对话轮次,而不是产品目录。在每条新消息中,agent 会加载当前会话的最近轮次,并将它们作为上下文重放,这样客户就不必重复自己说过的话。

这是集成式 Memory Management using AgentCore Memory lab 的自管理对应版本。它使用相同的能力和访问模式(每个会话的最近轮次),但构建在我们自己运行的基础设施之上。

Milvus 上的内存 vs RAG

我们现在已经将 Milvus 用于两种不同的任务。它们共享相同的存储引擎,但解决不同的问题:

RAG with Milvus Memory with Milvus(本实验)
存储的内容 产品目录 + FAQ embeddings(静态,共享) 客户的对话轮次(随时间增长)
索引键 无,一个共享目录 actor_id + session_id,用于该客户的会话
检索 向量搜索(“查找相似产品”) 时序性(“该会话的最近轮次,按顺序”)
目的 用产品知识为答案奠定基础 带着上下文继续对话
本实验按**会话内的时序性**进行检索。它在 Milvus 上使用标量 `query`(无向量搜索)。每一轮仍然**以 embedding 形式存储**,所以将其扩展为语义化的跨会话召回只需改一行代码:把 `query` 换成向量 `search`。

我们将构建的内容

一个客户服务 agent,在同一对话的多个轮次中:

  1. 从 Milvus 加载该 actor_id + session_id 的最近轮次
  2. 将它们重放到 agent 中,使得早先的上下文(例如订单 ID)得以延续
  3. 将每条新轮次记录回 Milvus,供下一条消息使用

工作原理

  1. 每一轮(客户问题 + agent 回答)都存储在 conversation_memory 集合中,并用 actor_idsession_id 和时间戳标记。
  2. 在新消息到来时,agent 运行一个按 actor_id + session_id 过滤、按时间排序的标量 query。结果是最近的轮次,最早的在前。consistency_level="Strong" 保证刚刚写入的轮次能立即可见。
  3. 这些轮次作为先前消息重放到 Strands agent 中。
  4. agent 带着该上下文回答,然后记录新的轮次。

步骤 1:验证 Milvus

kubectl get pods -n milvus

与 RAG 实验中相同的 milvus-standaloneetcdminio pods。内存复用它们,存储一个独立的 conversation_memory 集合。

步骤 2:代码讲解

cd ~/environment/modules/20-self-managed/500-memory-milvus/customer-agent

与 RAG 实验相比,rag_tools.pymemory.py 取代,而 agent.py 围绕 run_stream() 进行了重构,它会加载会话的最近轮次并记录新的轮次。其形态与集成式 AgentCore 实验相匹配,但由 Milvus 支撑。

@observe(name="milvus_memory.record_turn")
def record_turn(actor_id, session_id, user_message, assistant_message):
    """Persist one turn (embedded, so semantic recall stays possible later)."""
    _client.insert(COLLECTION, data=[{
        "actor_id": actor_id, "session_id": session_id,
        "user_message": user_message, "assistant_message": assistant_message,
        "ts": int(time.time()), "vector": _embed(user_message),
    }])

@observe(name="milvus_memory.recent_turns")
def recent_turns(actor_id, session_id, max_results=20):
    """This session's turns, oldest first. Scalar query, no vector search."""
    rows = _client.query(
        COLLECTION,
        filter=f'actor_id == "{actor_id}" && session_id == "{session_id}"',
        output_fields=["user_message", "assistant_message", "ts"],
        limit=max_results, consistency_level="Strong",
    )
    ...  # sort by ts, flatten to [{role, content}, ...]
  • 检索是对 actor_id + session_id标量 query。Milvus 充当过滤存储,而非向量搜索。
  • consistency_level="Strong" 是关键行。如果没有它,Milvus 默认的有界过期性可能会隐藏我们刚刚写入的轮次,从而破坏多轮召回。
  • 该轮次在写入时仍然被 embedded,所以该集合已准备好作为扩展用于语义召回。

def _build_agent(actor_id, session_id):
    prior = recent_turns(actor_id, session_id)   # this session's history
    return Agent(model=model, system_prompt=SYSTEM_PROMPT, tools=[lookup_order],
                 messages=[{"role": t["role"], "content": [{"text": t["content"]}]} for t in prior])

async def run_stream(actor_id, session_id, query):
    agent = _build_agent(actor_id, session_id)
    chunks = []
    async for event in agent.stream_async(query):
        if "data" in event:
            chunks.append(event["data"]); yield event["data"]
    record_turn(actor_id, session_id, query, "".join(chunks))  # persist for next turn

读取历史,用预加载的历史构建 agent,运行一轮,然后记录。与每个自管理实验一样,使用相同的 OpenAIModel 到 LiteLLM 到 vLLM 模型平面。只有内存层是新的。

步骤 3:部署 agent

envsubst < k8s.yaml | kubectl apply -f -
kubectl rollout status deployment/customer-agent --timeout=180s

步骤 4:跨轮次聊天

打开聊天 UI 并选择 Customer Agent (Self-managed GenAI)。UI 会为同一浏览器标签页中的每条消息发送一个稳定的 session_id,因此第 2 轮可以召回第 1 轮。

第 1 轮(建立上下文):

Hi, I'd like to check on my order ORD-12345.

第 2 轮(同一标签页,无订单 ID):

Has it shipped yet?

agent 会自行用 ORD-12345 调用 lookup_order,该 ID 是从第 1 轮存储的轮次中提取的,而不是我们重新输入的任何内容。内存层完成了这项工作。

打开一个新的浏览器标签页或重新加载 UI,我们将获得一个全新的 `session_id`,所以第 2 轮不会找到第 1 轮的上下文。这是会话作用域按设计工作,而不是 bug。

步骤 5:在 Langfuse 中查看

现在每一轮都会显示一个 milvus_memory.recent_turns span(读取)和一个 milvus_memory.record_turn span(写入),它们嵌套在 chat_turn 根节点之下。我们可以观察到内存在模型调用之前被读取,之后被写入。

我们构建的内容

Milvus 上的会话对话内存。agent 跨轮次携带上下文,而不是将每条消息都当作第一条来处理。与 RAG 实验相同的向量数据库,不同的集合,以及时序性访问模式。语义召回只差一行代码,因为轮次是以 embeddings 形式存储的。

下一步

工具仍然在进程内运行。接下来我们将在 Agent Tool Access (MCP) lab 中把它们移出到运行在 EKS 上的真正的 MCP server