实用说明:实验:Langgraph——集成Qdrant实现语义代理记忆功能。
《实用笔记:Lab:Langgraph——集成Qdrant实现语义代理记忆》的操作指南:为采用该模式的团队提供的契约、检查项以及可直接插入的代码片段。
可将此内容视为针对操作人员的《Lab:Langgraph: Integrating Qdrant for Semantic Agentic Memory.》中理念的重构版本:清晰的阶段划分、有序的代码模块,以及能在交接过程中保留的恢复说明。 将“概览”阶段视为可量化的界面使用效果最佳。在扩大范围之前,先记录一份理想的运行日志、一个故障案例以及回滚说明。 在功能结果旁同时记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。
为何需要语义记忆?
对于为何将语义记忆设为阶段这一设计,应在修改代码之前明确输入内容、该步骤的负责人以及退出标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储以及功能标志应集中存放于一个位置,以便操作人员无需查看整个流程即可进行审计。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不等同于业务流程的完整性。
fastapi>=0.110.0
uvicorn>=0.28.0
pydantic>=2.6.0
python-dotenv>=1.0.1
langchain-core>=0.1.30
langchain-openai>=0.1.0
langgraph>=0.0.30
qdrant-client>=1.10.0
#Install venv module (Ubuntu/Debian)
$sudo apt update && sudo apt install python3-venv
#Create a virtual environment
$python3 -m venv myenv
#Activate the environment
$source myenv/bin/activate
#Install saved dependencies
$pip install -r requirements.txt
#Save dependencies
#pip freeze > requirements.txt
#Deactivate when finished:
$deactivate
#Delete the environment:
$rm -rf myenv
#Install uv
$curl -LsSf [https://astral.sh/uv/install.sh](https://astral.sh/uv/install.sh) | sh
#Create a virtual environment
$uv venv myenv
#Activate the environment
$source myenv/bin/activate
#Install saved dependencies:
$uv pip install -r requirements.txt
#Save dependencies:
$uv pip freeze > requirements.txt
#Sync dependencies from lockfile
$uv sync
#Add a new package and auto-update file
$uv add <package_name>
version: '3.8'
services:
qdrant:
image: qdrant/qdrant:latest #image name download from docker.io
container_name: sbi_semantic_memory #container name
ports:
- "6333:6333" # REST HTTP API & Web Dashboard
- "6334:6334" # High-speed gRPC API
volumes:
- ./qdrant_storage:/qdrant/storage #Binds local directory to container for data persistence across restarts
networks:
- agent-network #Attaches container to isolated network for communication
healthcheck: #Executes internal HTTP ping against health check endpoint
test: ["CMD", "curl", "-f", "http://localhost:6333/healthz"]
interval: 10s #Run health check every 10 seconds
timeout: 5s #Fail if check takes longer than 5 seconds
retries: 5 #Startup grace period before recording failures
networks: #Private bridged network for secure container-to-container communication
agent-network:
driver: bridge #Standard single-host bridge driver
image: qdrant/qdrant:latest
$docker compose up -d
$ docker ps
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
8922a12f1029 qdrant/qdrant:latest "./entrypoint.sh" 45 hours ago Up 45 hours (unhealthy) 0.0.0.0:6333-6334->6333-6334/tcp, [::]:6333-6334->6333-6334/tcp sbi_semantic_memory
$ curl http://localhost:6333/healthz
healthz check passed
#View recent logs
$docker logs sbi_semantic_memory
#Tail live logs in real time
$docker logs -f sbi_semantic_memory
#View the last 50 lines of logs
$docker logs --tail 50 sbi_semantic_memory
#Inspect health check failure reasons
$docker inspect --format='{{json .State.Health}}' sbi_semantic_memory
#Check container resource consumption (CPU/RAM):
$docker stats sbi_semantic_memory
#Inspect container runtime details and exit codes:
$docker inspect sbi_semantic_memory
#Execute an interactive shell inside the container
$docker exec -it sbi_semantic_memory sh
#Restart the container:
$docker restart sbi_semantic_memory
#Stop, remove, and recreate the container:
$docker compose down && docker compose up -d
#Force-kill a stuck container:
$docker kill sbi_semantic_memory
import os
from dotenv import load_dotenv
from qdrant_client import QdrantClient
from qdrant_client.http import models
# Load environment variables
load_dotenv()
class SemanticMemory:
def __init__(self):
host = os.getenv("QDRANT_HOST", "localhost")
port = int(os.getenv("QDRANT_PORT", 6333))
self.client = QdrantClient(host=host, port=port)
self.collection_name = "agent_memories"
self._ensure_collection()
def _ensure_collection(self):
collections = self.client.get_collections().collections
exists = any(c.name == self.collection_name for c in collections)
if not exists:
self.client.create_collection(
collection_name=self.collection_name,
vectors_config=models.VectorParams(
size=1536,
distance=models.Distance.COSINE
),
)
memory_vault = SemanticMemory()
# OpenAI API Key
OPENAI_API_KEY=
# Qdrant Database Configuration
QDRANT_HOST=localhost
QDRANT_PORT=6333
QDRANT_COLLECTION_NAME=agent_memories
# FastAPI Server Setup
HOST=0.0.0.0
PORT=8000
from typing import Dict, Any, List
from dotenv import load_dotenv
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
# Ensure environment variables are loaded at application start
load_dotenv(override=True)
from graph import app_graph
app = FastAPI(title="Qdrant Semantic Memory Service")
class ChatRequest(BaseModel):
message: str
thread_id: str
metadata: Dict[str, Any] = {}
class ChatResponse(BaseModel):
thread_id: str
messages: List[Dict[str, Any]]
@app.post("/chat", response_model=ChatResponse)
async def chat_endpoint(req: ChatRequest):
try:
initial_state = {
"messages": [{"role": "user", "content": req.message}],
"thread_id": req.thread_id,
"metadata": req.metadata
}
final_state = await app_graph.ainvoke(initial_state)
return ChatResponse(
thread_id=req.thread_id,
messages=final_state["messages"]
)
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
{
"message": "My favorite color is Obsidian Blue.",
"thread_id": "thread-001",
"metadata": {"source": "mobile_app"}
}
{
"thread_id": "thread-001",
"messages": [
{
"role": "system",
"content": "You have access to the following long-term memory facts about the user:\n- My favorite color is Obsidian Blue.\n- My favorite color is Obsidian Blue.\n- Obsidian Blue is a deep, rich shade of blue that often resembles the color of the volcanic glass obsidian. It's a striking and elegant color choice! If you have any questions or need assistance related to colors or anything else, feel free to ask!\n\nUse the facts above to directly answer the user's question."
},
{
"role": "user",
"content": "My favorite color is Obsidian Blue."
},
{
"role": "assistant",
"content": "That's a great choice! Obsidian Blue is a deep, rich shade of blue that resembles the color of volcanic glass. It's both striking and elegant. If you have any questions or need assistance related to colors or anything else, feel free to ask!"
}
]
}
{
"message": "What is my favorite color?",
"thread_id": "thread-002",
"metadata": {"source": "web_dashboard" }
}
{
"thread_id": "thread-002",
"messages": [
{
"role": "system",
"content": "You have access to the following long-term memory facts about the user:\n- What is my favorite color?\n- My favorite color is Obsidian Blue.\n- My favorite color is Obsidian Blue.\n\nUse the facts above to directly answer the user's question."
},
{
"role": "user",
"content": "What is my favorite color?"
},
{
"role": "assistant",
"content": "Your favorite color is Obsidian Blue."
}
]
}
# 1. Generate query vector from user's message
query_vector = embeddings.embed_query(user_query)
# 2. Search Qdrant WITHOUT a payload filter
relevant_docs = memory_vault.client.query_points(
collection_name="agent_memories",
query=query_vector,
limit=2 # Grabs top 2 matches globally across all threads
).points
from qdrant_client.http import models
# Returns matches ONLY if thread_id matches current thread
relevant_docs = memory_vault.client.query_points(
collection_name="agent_memories",
query=query_vector,
query_filter=models.Filter(
must=[
models.FieldCondition(
key="thread_id",
match=models.MatchValue(value=state["thread_id"])
)
]
),
limit=2
).points
学到的核心概念:
在“关键概念学习”阶段,应在修改代码之前明确输入参数、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理措施都是产品本身的组成部分,而非后续需要补充的内容。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的配置并不等同于业务功能的完整性。
import os
import uuid
from typing import TypedDict, List, Dict, Any
from dotenv import load_dotenv
from langchain_openai import OpenAIEmbeddings, ChatOpenAI
from qdrant_client.http import models
from langgraph.graph import StateGraph, START, END
from app.memory.vector_store import memory_vault
load_dotenv()
class AgentState(TypedDict):
messages: List[Dict[str, Any]]
thread_id: str
tenant_id: str # Enforces multi-tenancy boundaries
user_id: str # Identifies the end-user
metadata: Dict[str, Any]
embeddings = OpenAIEmbeddings(model="text-embedding-3-small")
llm = ChatOpenAI(model="gpt-4o", temperature=0)
async def recall_node(state: AgentState) -> dict:
"""Retrieves semantic memories STRICTLY bounded by tenant_id and user_id."""
user_query = state["messages"][-1]["content"]
query_vector = embeddings.embed_query(user_query)
# Multi-Tenant Payload Filter Enforcement
tenant_filter = models.Filter(
must=[
models.FieldCondition(
key="tenant_id",
match=models.MatchValue(value=state["tenant_id"])
),
models.FieldCondition(
key="user_id",
match=models.MatchValue(value=state["user_id"])
)
]
)
# Scoped Query Execution
relevant_docs = memory_vault.client.query_points(
collection_name="agent_memories",
query=query_vector,
query_filter=tenant_filter, # Prevents cross-tenant data leaks
limit=3
).points
if relevant_docs:
memory_context = "\n".join([f"- {hit.payload['text']}" for hit in relevant_docs])
system_instruction = {
"role": "system",
"content": (
"You have access to the following long-term memory facts about the user:\n"
f"{memory_context}\n\n"
"Use the facts above to directly answer the user's question."
)
}
state["messages"].insert(0, system_instruction)
return {"messages": state["messages"]}
async def agent_node(state: AgentState) -> dict:
"""LLM reasoning node."""
response = await llm.ainvoke(state["messages"])
state["messages"].append({"role": "assistant", "content": response.content})
return {"messages": state["messages"]}
async def memorize_node(state: AgentState) -> dict:
"""Embeds user facts along with mandatory multi-tenant metadata payloads."""
user_message = next(
(m["content"] for m in reversed(state["messages"]) if m.get("role") == "user"),
None
)
if user_message:
vector = embeddings.embed_query(user_message)
memory_vault.client.upsert(
collection_name="agent_memories",
points=[
models.PointStruct(
id=str(uuid.uuid4()),
vector=vector,
payload={
"text": user_message,
"tenant_id": state["tenant_id"], # Tenant payload tag
"user_id": state["user_id"], # User payload tag
"thread_id": state["thread_id"],
"metadata": state.get("metadata", {})
}
)
]
)
return state
# Graph Workflow Wiring
workflow = StateGraph(AgentState)
workflow.add_node("recall", recall_node)
workflow.add_node("agent", agent_node)
workflow.add_node("memorize", memorize_node)
workflow.add_edge(START, "recall")
workflow.add_edge("recall", "agent")
workflow.add_edge("agent", "memorize")
workflow.add_edge("memorize", END)
app_graph = workflow.compile()
操作检查清单
在完成“操作检查清单”阶段时,首先需明确相关约定:所需的输入参数、成功信号,以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改保持一致性。
将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,明确成功标准,并拒绝默许的半完成状态。
在成本较高的步骤之后设置检查点。当操作员重新尝试后续节点时,恢复流程不应再次收取相同的LLM调用费用。
锁定依赖项的版本,并记录用于运行演示的镜像摘要。可重复性比经验知识更为可靠。
在功能结果旁记录执行时间以及token或查询成本。提前了解成本可避免从演示环境过渡到共享环境时出现意外收费。
在成本较高的步骤之后设置检查点。当操作员重新尝试后续节点时,恢复流程不应再次收取相同的LLM调用费用。
在推广该技术栈之前,应先冻结版本,为关键路径生成标准记录,并确认回滚步骤。共享环境需要设置速率限制、租户验证机制,以及明确的密钥轮换负责人。与其展示花哨的一次性演示,不如追求扎实的可靠性。
e7110284d0c4的批量处理说明:不要将提供商密钥放入代码仓库,为每个会话设置令牌上限,并将记录存储在评估用示例文件旁边,以便后续模型更换时仍能保持可比性。
针对强化安全性的第0阶段,在修改代码之前需明确输入内容、该步骤的负责人以及完成标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。应将此阶段视为输入与验证后输出之间的契约,为相关文件命名、定义成功检测标准,并拒绝默许的半完成状态。
安全加固细节 0/759:为该记录测量运行时间、错误类型以及代币消耗情况,然后依据固定的问题清单而非个人经验来判断是否保留该变更。
在处理安全加固笔记的第一阶段时,首先写下合约的必要输入参数、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改更加规范。 应将配置信息置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。
安全加固细节 1/759:为该记录测量运行时间、错误类型以及代币消耗情况,然后依据固定的问题清单而非个人经验来判断是否保留该变更。
在将加固措施视为可测量的表面时,第二阶段的效果最佳。在扩大范围之前,先记录一份理想的测试结果、一个故障案例以及回滚说明。 相比庞大的脚本,应优先选择小型且可测试的单元。当某个步骤出现故障时,故障原因应能指向单一责任方,而非复杂的流程链。
加固细节 2/759:针对此条说明,需测量执行时间、错误类型以及令牌消耗情况,然后依据固定的问题集而非个人经验来判断是否保留该变更。