title: “使用 Milvus 进行内存管理” weight: 500
在本模块中,我们将使用与 RAG with Milvus lab 相同的 Milvus 向量数据库为 agent 提供对话内存。这一次我们存储的是客户自己的对话轮次,而不是产品目录。在每条新消息中,agent 会加载当前会话的最近轮次,并将它们作为上下文重放,这样客户就不必重复自己说过的话。
这是集成式 Memory Management using AgentCore Memory lab 的自管理对应版本。它使用相同的能力和访问模式(每个会话的最近轮次),但构建在我们自己运行的基础设施之上。
我们现在已经将 Milvus 用于两种不同的任务。它们共享相同的存储引擎,但解决不同的问题:
| RAG with Milvus | Memory with Milvus(本实验) | |
|---|---|---|
| 存储的内容 | 产品目录 + FAQ embeddings(静态,共享) | 客户的对话轮次(随时间增长) |
| 索引键 | 无,一个共享目录 | actor_id + session_id,用于该客户的会话 |
| 检索 | 向量搜索(“查找相似产品”) | 时序性(“该会话的最近轮次,按顺序”) |
| 目的 | 用产品知识为答案奠定基础 | 带着上下文继续对话 |
本实验按**会话内的时序性**进行检索。它在 Milvus 上使用标量 `query`(无向量搜索)。每一轮仍然**以 embedding 形式存储**,所以将其扩展为语义化的跨会话召回只需改一行代码:把 `query` 换成向量 `search`。
一个客户服务 agent,在同一对话的多个轮次中:
actor_id + session_id 的最近轮次conversation_memory 集合中,并用 actor_id、session_id 和时间戳标记。actor_id + session_id 过滤、按时间排序的标量 query。结果是最近的轮次,最早的在前。consistency_level="Strong" 保证刚刚写入的轮次能立即可见。
kubectl get pods -n milvus
与 RAG 实验中相同的 milvus-standalone、etcd 和 minio pods。内存复用它们,存储一个独立的 conversation_memory 集合。
cd ~/environment/modules/20-self-managed/500-memory-milvus/customer-agent
与 RAG 实验相比,rag_tools.py 被 memory.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 默认的有界过期性可能会隐藏我们刚刚写入的轮次,从而破坏多轮召回。
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 模型平面。只有内存层是新的。
envsubst < k8s.yaml | kubectl apply -f -
kubectl rollout status deployment/customer-agent --timeout=180s
打开聊天 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。
现在每一轮都会显示一个 milvus_memory.recent_turns span(读取)和一个 milvus_memory.record_turn span(写入),它们嵌套在 chat_turn 根节点之下。我们可以观察到内存在模型调用之前被读取,之后被写入。
Milvus 上的会话对话内存。agent 跨轮次携带上下文,而不是将每条消息都当作第一条来处理。与 RAG 实验相同的向量数据库,不同的集合,以及时序性访问模式。语义召回只差一行代码,因为轮次是以 embeddings 形式存储的。
工具仍然在进程内运行。接下来我们将在 Agent Tool Access (MCP) lab 中把它们移出到运行在 EKS 上的真正的 MCP server。