Практические заметки: Лабораторная работа Langgraph: интеграция Qdrant для семантической агентной памяти.
Пошаговое руководство по использованию практических заметок: лабораторная работа 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()
Чек-лист операций
При работе над чек-листом операций сначала необходимо описать условия работы: требуемые входные данные, сигнал о успехе и действия при частичной неудаче. Такой чек-лист помогает сохранять честность при последующих изменениях кода.
Рассматривайте этот этап как контракт между входными данными и проверенными выходными результатами. Дайте названия результатам работы, определите критерии успешного выполнения и не допускайте молчаливого частичного завершения задачи.
Создавайте контрольные точки после дорогостоящих операций. Функция возобновления работы не должна снова взимать плату за один и тот же вызов большой языковой модели при повторной попытке обработки последующего элемента.
Закрепите версии зависимостей и запишите хэш изображения, с использованием которого выполнялась демонстрация. Воспроизводимость важнее устного опыта специалистов.
Записывайте время выполнения операций, а также стоимость токенов или запросов рядом с функциональными результатами. Ясность стоимости на ранних этапах предотвращает неожиданные счета при переходе от демо-среды к общедоступным средам.
Создавайте контрольные точки после дорогостоящих операций. Функция возобновления работы не должна снова взимать плату за один и тот же вызов большой языковой модели при повторной попытке обработки последующего элемента.
Перед внедрением стека заморозьте версии, сделайте копию «золотого» транскрипта для критической цепочки операций и уточните шаги отката. В совместных средах необходимы ограничения по частоте запросов, проверки принадлежности пользователя и четко определенный ответственный за обновление секретов. Лучше выбирать надежность, даже если она кажется скучной, чем умные одноразовые демонстрации.
Примечание для e7110284d0c4: не храните ключи поставщика в репозитории, установите лимит токенов на одну сессию и сохраняйте транскрипты рядом с фикстурами для оценки, чтобы последующие замены моделей оставались сопоставимыми.
Для этапа 0 по усилению безопасности определите входные данные, ответственного за шаг и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Рассматривайте этот этап как контракт между входными данными и проверенными выходными результатами. Дайте названия артефактам, определите критерии успеха и не допускайте безответственного частичного выполнения задачи.
Подробности усиления безопасности 0/759: измерьте время выполнения, класс ошибки и расход токенов для этой записи, затем решите, следует ли сохранить изменение на основе фиксированного набора вопросов, а не на основе единичных примеров.
При работе над первым этапом записи по усилению безопасности сначала запишите контракт: необходимые входные данные, сигнал успешного выполнения и то, что происходит при частичной неудаче. Такой чек-лист помогает сохранять честность последующих изменений в коде. Храните конфигурацию вне кода приложения. Файлы среды, хранилища секретов и флаги функций должны находиться в одном месте, которое операторы могут проверять, не читая весь код.
Подробности усиления безопасности 1/759: измерьте время выполнения, класс ошибки и расход токенов для этой записи, затем решите, следует ли сохранить изменение на основе фиксированного набора вопросов, а не на основе единичных примеров.
Второй этап усиления безопасности работает наилучшим образом, когда его рассматривают как измеримую поверхность. Соберите один идеальный пример работы, один случай сбоя и запись о возврате к предыдущему состоянию перед расширением объема работ. Предпочитайте небольшие, тестируемые единицы кода вместо обширных скриптов. Когда какой-либо шаг терпит неудачу, причина сбоя должна указывать на конкретный ответственный элемент, а не на запутанную цепочку операций.
Подробности усиления безопасности 2/759: измерьте время выполнения, класс ошибки и количество использованных токенов для этой записи, затем решите, следует ли сохранять изменения, опираясь на заранее определенный набор критериев, а не на устные описания.