第3章 存储技术实现 本章将深入探讨记忆系统的存储技术实现,从数据库选型到具体的存储策略,为读者提供详细的存储解决方案。 3.1 存储技术选型 3.1.1 存储需求分析 在开始存储技术选型之前,我们需要明确记忆系统的存储需求: 数据类型多样:需要支持文本、图像、音频、向量等多种数据类型 查询模式复杂:支持关键词搜索、语义搜索、混合查询等多种查询方式 性能要求高:需要毫秒级的响应时间,特别是在读取操作上 数据一致性:保证数据的一致性和完整性,特别是在分布式环境下 可扩展性:支持水平扩展,能够应对数据量和访问量的增长 成本控制:在性能和成本之间找到平衡点 3.1.
本章将深入探讨记忆系统的存储技术实现,从数据库选型到具体的存储策略,为读者提供详细的存储解决方案。
在开始存储技术选型之前,我们需要明确记忆系统的存储需求:
| 存储技术 | 适用场景 | 优点 | 缺点 | 成本 |
|---|---|---|---|---|
| MySQL | 结构化数据,关系查询 | 成熟稳定,事务支持好 | 扩展性差,文本搜索弱 | 中 |
| PostgreSQL | 复杂查询,全文搜索 | 功能强大,扩展性好 | 学习成本高 | 中高 |
| MongoDB | 文档存储,灵活模式 | 灵活扩展,开发效率高 | 一致性较弱,性能波动 | 中 |
| Redis | 高性能缓存,内存存储 | 极快读写,丰富数据结构 | 内存成本高,数据量有限 | 高 |
| Elasticsearch | 全文搜索,日志分析 | 搜索性能优异,扩展性好 | 资源消耗大,配置复杂 | 高 |
| Milvus | 向量存储,相似度搜索 | 专门优化向量搜索,高性能 | 专业性强,生态较小 | 中高 |
| Qdrant | 轻量级向量数据库 | 轻量级,易部署 | 功能相对简单 | 中 |
| Pinecone | 托管向量数据库 | 免运维,扩展性好 | 成本较高,灵活性受限 | 高 |
基于上述分析,我们推荐采用混合存储架构:
MySQL作为主要的关系型数据库,用于存储记忆的基础信息:
-- 记忆条目主表 CREATE TABLE memory_items ( id VARCHAR(36) PRIMARY KEY, type ENUM('fact', 'procedure', 'contextual', 'emotional') NOT NULL, category VARCHAR(50) NOT NULL, title VARCHAR(200) NOT NULL, content TEXT NOT NULL, content_type VARCHAR(50) NOT NULL, metadata JSON, access_count INT DEFAULT 0, last_accessed TIMESTAMP NULL, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, version INT DEFAULT 1, is_active BOOLEAN DEFAULT TRUE, importance DECIMAL(3,2) DEFAULT 0.5, expiry_time TIMESTAMP NULL, INDEX idx_type_category (type, category), INDEX idx_created_at (created_at), INDEX idx_importance (importance), INDEX idx_expiry_time (expiry_time) ) ENGINE=InnoDB; -- 记忆标签表 CREATE TABLE memory_tags ( id VARCHAR(36) PRIMARY KEY, memory_id VARCHAR(36) NOT NULL, tag VARCHAR(100) NOT NULL, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, FOREIGN KEY (memory_id) REFERENCES memory_items(id) ON DELETE CASCADE, INDEX idx_memory_id (memory_id), INDEX idx_tag (tag), UNIQUE KEY unique_memory_tag (memory_id, tag) ) ENGINE=InnoDB; -- 记忆访问日志表 CREATE TABLE memory_access_logs ( id VARCHAR(36) PRIMARY KEY, memory_id VARCHAR(36) NOT NULL, access_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, access_type ENUM('read', 'write', 'delete') NOT NULL, access_context JSON, user_ip VARCHAR(45), user_agent VARCHAR(255), FOREIGN KEY (memory_id) REFERENCES memory_items(id) ON DELETE CASCADE, INDEX idx_memory_id (memory_id), INDEX idx_access_time (access_time) ) ENGINE=InnoDB; -- 记忆版本表 CREATE TABLE memory_versions ( id VARCHAR(36) PRIMARY KEY, memory_id VARCHAR(36) NOT NULL, version_number INT NOT NULL, content TEXT NOT NULL, metadata JSON, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, created_by VARCHAR(36), change_reason VARCHAR(500), FOREIGN KEY (memory_id) REFERENCES memory_items(id) ON DELETE CASCADE, INDEX idx_memory_id (memory_id), INDEX idx_version_number (version_number) ) ENGINE=InnoDB;
使用Python实现数据操作:
import mysql.connector from mysql.connector import Error from datetime import datetime, timedelta import json import uuid class MemoryMySQLRepository: def __init__(self, config): self.config = config self.connection = None def connect(self): """建立数据库连接""" try: self.connection = mysql.connector.connect(**self.config) return True except Error as e: print(f"数据库连接失败: {e}") return False def create_memory(self, memory_data): """创建记忆条目""" if not self.connection: self.connect() try: cursor = self.connection.cursor() # 生成唯一ID memory_id = str(uuid.uuid4()) # 主插入语句 query = """ INSERT INTO memory_items ( id, type, category, title, content, content_type, metadata, importance, expiry_time ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) """ values = ( memory_id, memory_data['type'], memory_data['category'], memory_data['title'], memory_data['content'], memory_data['content_type'], json.dumps(memory_data.get('metadata', {})), memory_data.get('importance', 0.5), memory_data.get('expiry_time') ) cursor.execute(query, values) # 插入标签 if 'tags' in memory_data: tag_query = """ INSERT INTO memory_tags (id, memory_id, tag) VALUES (%s, %s, %s) """ for tag in memory_data['tags']: tag_id = str(uuid.uuid4()) cursor.execute(tag_query, (tag_id, memory_id, tag)) self.connection.commit() return memory_id except Error as e: self.connection.rollback() raise e finally: if cursor: cursor.close() def get_memory(self, memory_id): """获取记忆条目""" if not self.connection: self.connect() try: cursor = self.connection.cursor(dictionary=True) query = """ SELECT * FROM memory_items WHERE id = %s AND is_active = TRUE """ cursor.execute(query, (memory_id,)) result = cursor.fetchone() if result: # 获取标签 cursor.execute(""" SELECT tag FROM memory_tags WHERE memory_id = %s """, (memory_id,)) result['tags'] = [row['tag'] for row in cursor.fetchall()] # 更新访问统计 self.update_access_count(memory_id) return result except Error as e: raise e finally: if cursor: cursor.close() def search_memories(self, query_params): """搜索记忆条目""" if not self.connection: self.connect() try: cursor = self.connection.cursor(dictionary=True) # 构建基础查询 base_query = """ SELECT * FROM memory_items WHERE is_active = TRUE """ conditions = [] params = [] # 类型过滤 if 'type' in query_params: conditions.append("type = %s") params.append(query_params['type']) # 分类过滤 if 'category' in query_params: conditions.append("category = %s") params.append(query_params['category']) # 标签过滤 if 'tags' in query_params: tag_conditions = [] for tag in query_params['tags']: tag_conditions.append("id IN (SELECT memory_id FROM memory_tags WHERE tag = %s)") params.append(tag) if tag_conditions: conditions.append("(" + " OR ".join(tag_conditions) + ")") # 构建完整查询 if conditions: base_query += " AND " + " AND ".join(conditions) # 排序 if 'sort_by' in query_params: base_query += f" ORDER BY {query_params['sort_by']} {query_params.get('sort_order', 'DESC')}" else: base_query += " ORDER BY last_accessed DESC, importance DESC" # 分页 if 'limit' in query_params: base_query += " LIMIT %s" params.append(query_params['limit']) if 'offset' in query_params: base_query += " OFFSET %s" params.append(query_params['offset']) cursor.execute(base_query, params) results = cursor.fetchall() # 为每个结果获取标签 for result in results: cursor.execute(""" SELECT tag FROM memory_tags WHERE memory_id = %s """, (result['id'],)) result['tags'] = [row['tag'] for row in cursor.fetchall()] return results except Error as e: raise e finally: if cursor: cursor.close()
MySQL的性能优化对于记忆系统至关重要:
# MySQL性能优化配置 class MySQLPerformanceConfig: def __init__(self): # 连接池配置 self.connection_pool = { 'pool_name': 'memory_pool', 'pool_size': 20, 'pool_reset_session': True, 'pool_recycle': 3600, 'pool_pre_ping': True } # 索引优化 self.index_config = { 'composite_indexes': [ ('type', 'category', 'is_active'), ('importance', 'last_accessed'), ('created_at', 'expiry_time') ], 'fulltext_indexes': [ ('content',), ('title',) ] } # 查询优化 self.query_config = { 'use_read_replicas': True, 'query_cache': True, 'slow_query_threshold': 1000, # 1秒 'large_query_threshold': 10000 # 10000行 } # 分区策略 self.partition_config = { 'enabled': True, 'strategy': 'RANGE', 'column': 'created_at', 'interval': '1 MONTH', 'retention': '12 MONTH' }
Qdrant是一个轻量级的向量数据库,非常适合记忆系统的向量存储需求:
import qdrant_client from qdrant_client.http import models from qdrant_client.http.models import Distance, VectorParams import numpy as np from typing import List, Dict, Optional class MemoryVectorRepository: def __init__(self, url: str, api_key: str = None): self.client = qdrant_client.QdrantClient( url=url, api_key=api_key ) self.collection_name = "memory_vectors" # 确保集合存在 self._ensure_collection() def _ensure_collection(self): """确保向量集合存在""" try: self.client.get_collection(self.collection_name) except: # 集合不存在,创建新的 self.client.create_collection( collection_name=self.collection_name, vectors_config=VectorParams( size=1536, # OpenAI embedding维度 distance=Distance.COSINE ), optimizers_config=models.OptimizersConfigDiff( indexing_threshold=20000 ), hnsw_config=models.HnswConfigDiff( m=16, ef=64 ) ) def add_memory_vector(self, memory_id: str, text: str, embedding: List[float], metadata: Dict = None): """添加记忆向量""" point_id = f"memory_{memory_id}" # 准备点数据 point = models.PointStruct( id=point_id, vector=embedding, payload={ "memory_id": memory_id, "text": text, "metadata": metadata or {}, "created_at": datetime.now().isoformat(), "type": "memory_vector" } ) # 批量添加点 self.client.upsert( collection_name=self.collection_name, points=[point] ) def search_memories(self, query_embedding: List[float], limit: int = 10, score_threshold: float = 0.7) -> List[Dict]: """搜索相似记忆""" search_result = self.client.search( collection_name=self.collection_name, query_vector=query_embedding, limit=limit, score_threshold=score_threshold, with_payload=True, search_filter=models.Filter( must=[ models.FieldCondition( key="type", match=models.MatchValue(value="memory_vector") ) ] ) ) # 转换结果 results = [] for point in search_result: results.append({ "memory_id": point.payload["memory_id"], "text": point.payload["text"], "score": point.score, "metadata": point.payload["metadata"], "created_at": point.payload["created_at"] }) return results def delete_memory_vector(self, memory_id: str): """删除记忆向量""" point_id = f"memory_{memory_id}" self.client.delete( collection_name=self.collection_name, points_selector=models.PointIdsList( points=[point_id] ) )
Redis作为高性能缓存,为记忆系统提供快速的访问能力:
import redis import json from typing import Optional, Dict, Any import pickle class MemoryCacheService: def __init__(self, config: Dict): self.redis_client = redis.Redis( host=config['host'], port=config['port'], db=config.get('db', 0), password=config.get('password'), decode_responses=True ) # 本地缓存 self.local_cache = {} self.local_cache_size = 1000 # 缓存键前缀 self.prefix = config.get('prefix', 'memory:') # 缓存过期时间 self.ttl = config.get('ttl', 3600) # 1小时 def get_memory(self, memory_id: str) -> Optional[Dict]: """获取记忆缓存""" # 1. 先查本地缓存 cache_key = f"{self.prefix}memory:{memory_id}" if cache_key in self.local_cache: return self._deserialize(self.local_cache[cache_key]) # 2. 查Redis缓存 try: cached_data = self.redis_client.get(cache_key) if cached_data: data = self._deserialize(cached_data) # 更新本地缓存 self._update_local_cache(cache_key, cached_data) return data except Exception as e: print(f"Redis缓存查询失败: {e}") return None def set_memory(self, memory_id: str, data: Dict, ttl: int = None): """设置记忆缓存""" cache_key = f"{self.prefix}memory:{memory_id}" # 序列化数据 serialized_data = self._serialize(data) # 设置Redis缓存 try: redis_ttl = ttl or self.ttl self.redis_client.setex(cache_key, redis_ttl, serialized_data) # 更新本地缓存 self._update_local_cache(cache_key, serialized_data) except Exception as e: print(f"Redis缓存设置失败: {e}") def batch_get_memories(self, memory_ids: List[str]) -> Dict[str, Dict]: """批量获取记忆缓存""" results = {} missing_ids = [] # 先查询本地缓存 for memory_id in memory_ids: cache_key = f"{self.prefix}memory:{memory_id}" if cache_key in self.local_cache: results[memory_id] = self._deserialize(self.local_cache[cache_key]) else: missing_ids.append(memory_id) # 查询Redis缓存 if missing_ids: try: redis_keys = [f"{self.prefix}memory:{mid}" for mid in missing_ids] cached_data = self.redis_client.mget(redis_keys) for i, data in enumerate(cached_data): if data: memory_id = missing_ids[i] results[memory_id] = self._deserialize(data) self._update_local_cache(redis_keys[i], data) except Exception as e: print(f"Redis批量缓存查询失败: {e}") return results def _serialize(self, data: Any) -> str: """序列化数据""" return json.dumps(data, ensure_ascii=False) def _deserialize(self, data: str) -> Any: """反序列化数据""" return json.loads(data) def _update_local_cache(self, key: str, data: str): """更新本地缓存""" # 如果本地缓存已满,删除最旧的项 if len(self.local_cache) >= self.local_cache_size: oldest_key = next(iter(self.local_cache)) del self.local_cache[oldest_key] self.local_cache[key] = data
本章详细介绍了记忆系统的存储技术实现,包括关系型数据库、向量数据库和缓存系统的具体实现方案。通过合理的技术选型和配置,可以构建一个高效、可扩展的记忆存储系统。
在下一章中,我们将探讨记忆系统的检索引擎实现,包括关键词搜索、语义搜索和混合检索等技术方案。