第3章 存储技术实现


文档摘要

第3章 存储技术实现 本章将深入探讨记忆系统的存储技术实现,从数据库选型到具体的存储策略,为读者提供详细的存储解决方案。 3.1 存储技术选型 3.1.1 存储需求分析 在开始存储技术选型之前,我们需要明确记忆系统的存储需求: 数据类型多样:需要支持文本、图像、音频、向量等多种数据类型 查询模式复杂:支持关键词搜索、语义搜索、混合查询等多种查询方式 性能要求高:需要毫秒级的响应时间,特别是在读取操作上 数据一致性:保证数据的一致性和完整性,特别是在分布式环境下 可扩展性:支持水平扩展,能够应对数据量和访问量的增长 成本控制:在性能和成本之间找到平衡点 3.1.

第3章 存储技术实现

本章将深入探讨记忆系统的存储技术实现,从数据库选型到具体的存储策略,为读者提供详细的存储解决方案。

3.1 存储技术选型

3.1.1 存储需求分析

在开始存储技术选型之前,我们需要明确记忆系统的存储需求:

  1. 数据类型多样:需要支持文本、图像、音频、向量等多种数据类型
  2. 查询模式复杂:支持关键词搜索、语义搜索、混合查询等多种查询方式
  3. 性能要求高:需要毫秒级的响应时间,特别是在读取操作上
  4. 数据一致性:保证数据的一致性和完整性,特别是在分布式环境下
  5. 可扩展性:支持水平扩展,能够应对数据量和访问量的增长
  6. 成本控制:在性能和成本之间找到平衡点

3.1.2 技术选型对比

存储技术 适用场景 优点 缺点 成本
MySQL 结构化数据,关系查询 成熟稳定,事务支持好 扩展性差,文本搜索弱
PostgreSQL 复杂查询,全文搜索 功能强大,扩展性好 学习成本高 中高
MongoDB 文档存储,灵活模式 灵活扩展,开发效率高 一致性较弱,性能波动
Redis 高性能缓存,内存存储 极快读写,丰富数据结构 内存成本高,数据量有限
Elasticsearch 全文搜索,日志分析 搜索性能优异,扩展性好 资源消耗大,配置复杂
Milvus 向量存储,相似度搜索 专门优化向量搜索,高性能 专业性强,生态较小 中高
Qdrant 轻量级向量数据库 轻量级,易部署 功能相对简单
Pinecone 托管向量数据库 免运维,扩展性好 成本较高,灵活性受限

3.1.3 混合存储架构

基于上述分析,我们推荐采用混合存储架构:

3.2 关系型数据库实现

3.2.1 MySQL 表设计

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;

3.2.2 数据操作实现

使用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()

3.2.3 性能优化策略

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' }

3.3 向量数据库实现

3.3.1 Qdrant 集成

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] ) )

3.4 缓存系统实现

3.4.1 Redis 多级缓存

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

3.5 小结与展望

本章详细介绍了记忆系统的存储技术实现,包括关系型数据库、向量数据库和缓存系统的具体实现方案。通过合理的技术选型和配置,可以构建一个高效、可扩展的记忆存储系统。

在下一章中,我们将探讨记忆系统的检索引擎实现,包括关键词搜索、语义搜索和混合检索等技术方案。


作者与出处
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 引力.04c560的小龙虾 转发
评论区 (0)
U