4.1 系统架构设计 本节导读:本节将从工程实现的角度,详细展示 GraphRAG 系统的五层架构设计。我们将从架构概览出发,逐一讲解数据接入层、知识图谱层、检索引擎层、应用服务层和存储层的核心组件设计、接口定义和协作机制。通过完整的类设计和代码示例,帮助读者建立一个可以在生产环境中直接使用的 GraphRAG 系统骨架。 学习目标 掌握 GraphRAG 系统的五层架构设计原则 理解各层核心组件的职责与接口定义 掌握查询处理引擎的完整实现 学会知识图谱管理器的 CRUD 操作封装 理解容器化部署与编排策略 架构设计原则 为什么需要分层架构 GraphRAG 系统的实现是一个典型的分层工程问题。
本节导读:本节将从工程实现的角度,详细展示 GraphRAG 系统的五层架构设计。我们将从架构概览出发,逐一讲解数据接入层、知识图谱层、检索引擎层、应用服务层和存储层的核心组件设计、接口定义和协作机制。通过完整的类设计和代码示例,帮助读者建立一个可以在生产环境中直接使用的 GraphRAG 系统骨架。
GraphRAG 系统的实现是一个典型的分层工程问题。与传统的 RAG 系统不同,GraphRAG 在文本检索之外引入了知识图谱的结构化推理能力,这使得系统架构的复杂度显著提升。
在单一进程中实现所有功能虽然简单,但在生产环境中会面临严峻挑战:
GraphRAG 系统的五层架构如下:
数据接入层是系统的入口——负责从各种数据源(PDF、网页、数据库、API)中摄取原始文本,并进行初步的清洗和分块处理。核心挑战在于处理异构数据源的格式差异,确保进入系统的文本质量。
知识图谱层是系统的核心——负责将清洗后的文本转化为结构化的实体-关系三元组,并维护知识图谱的持续更新。集成了 NER 模型、关系抽取器和图数据库,是系统中最复杂的部分。
检索引擎层是系统的查询处理器——负责接收用户查询,在知识图谱和向量索引中执行多路检索,并通过混合排序算法生成最终结果。直接决定了系统的检索质量和响应延迟。
应用服务层是系统的对外接口——负责暴露 RESTful API,处理认证、限流、日志等横切关注点,并将检索结果格式化为前端友好的结构。
存储层是系统的持久化基础——提供关系数据存储、高性能缓存和二进制对象存储能力。
五层之间的协作遵循单向依赖原则:上层调用下层提供的接口,但下层不感知上层的存在。依赖关系通过抽象接口实现。
查询处理的数据流路径:
这种并行查询 + 融合排序的架构是 GraphRAG 区别于传统 RAG 的关键特征。
查询处理引擎是检索引擎层的核心组件,负责从接收查询到生成检索结果的完整流程。
from dataclasses import dataclass, field from typing import List, Dict, Optional, Protocol from enum import Enum import asyncio class QueryType(Enum): """查询类型枚举""" ENTITY = "entity" # 实体查询:"张三属于哪个部门?" RELATION = "relation" # 关系查询:"A系统和B系统有什么关系?" DESCRIPTIVE = "descriptive" # 描述性查询:"什么是GraphRAG?" COMPLEX = "complex" # 复合查询:"张三的部门负责的项目用了哪些技术?" @dataclass class QueryResult: """统一检索结果""" content: str source: str # "graph" | "vector" | "merged" score: float metadata: Dict = field(default_factory=dict) path: Optional[List[str]] = None # 知识图谱路径(图检索时) class KnowledgeGraphClient(Protocol): """知识图谱客户端抽象接口""" async def search_paths(self, entity: str, max_hops: int = 3) -> List[QueryResult]: ... async def get_neighbors(self, entity: str) -> List[Dict]: ... async def search_by_pattern(self, pattern: Dict) -> List[QueryResult]: ... class VectorStoreClient(Protocol): """向量存储客户端抽象接口""" async def search(self, query_vector: List[float], top_k: int = 10) -> List[QueryResult]: ... class QueryProcessor: """查询处理引擎""" def __init__(self, kg_client: KnowledgeGraphClient, vector_client: VectorStoreClient, graph_weight: float = 0.6): self.kg_client = kg_client self.vector_client = vector_client self.graph_weight = graph_weight self.vector_weight = 1.0 - graph_weight def classify_query(self, query: str) -> QueryType: """查询类型分类""" # 简化的分类逻辑:基于关键词和实体数量 entity_count = self._extract_entities(query) if "关系" in query or "联系" in query or "关联" in query: return QueryType.RELATION if entity_count >= 2: return QueryType.COMPLEX if entity_count == 1: return QueryType.ENTITY return QueryType.DESCRIPTIVE async def process(self, query: str, top_k: int = 10) -> List[QueryResult]: """处理查询请求""" query_type = self.classify_query(query) # 根据查询类型调整权重 if query_type in (QueryType.ENTITY, QueryType.RELATION, QueryType.COMPLEX): graph_weight = min(self.graph_weight + 0.2, 0.8) else: graph_weight = max(self.graph_weight - 0.2, 0.3) # 并行执行两路检索 graph_task = asyncio.create_task( self._graph_search(query, top_k) ) vector_task = asyncio.create_task( self._vector_search(query, top_k) ) graph_results, vector_results = await asyncio.gather( graph_task, vector_task ) # 融合排序 merged = self._reciprocal_rank_fusion( graph_results, vector_results, graph_weight, 1.0 - graph_weight ) return merged[:top_k] def _extract_entities(self, query: str) -> int: """提取查询中的实体数量(简化版)""" # 实际项目中应使用 NER 模型 return sum(1 for c in query if c.isupper() and c.isalpha()) async def _graph_search(self, query: str, top_k: int) -> List[QueryResult]: """图结构检索""" entities = self._extract_entities(query) if not entities: return [] results = [] for entity in entities[:3]: paths = await self.kg_client.search_paths(entity, max_hops=2) results.extend(paths) return results async def _vector_search(self, query: str, top_k: int) -> List[QueryResult]: """向量语义检索""" query_vector = self._embed(query) results = await self.vector_client.search(query_vector, top_k) return results def _embed(self, text: str) -> List[float]: """文本向量化(需替换为实际的嵌入模型)""" # placeholder return [0.0] * 768 def _reciprocal_rank_fusion(self, results_a: List[QueryResult], results_b: List[QueryResult], weight_a: float, weight_b: float, k: int = 60) -> List[QueryResult]: """倒数排名融合""" scores = {} for rank, item in enumerate(results_a): key = item.content scores[key] = scores.get(key, 0) + weight_a / (k + rank + 1) for rank, item in enumerate(results_b): key = item.content scores[key] = scores.get(key, 0) + weight_b / (k + rank + 1) merged = sorted(scores.items(), key=lambda x: x[1], reverse=True) return [QueryResult(content=item[0], score=item[1], source="merged") for item in merged]
知识图谱管理器负责知识图谱的 CRUD 操作,封装图数据库的底层操作。
from typing import List, Dict, Optional from dataclasses import dataclass @dataclass class Entity: """知识图谱实体""" name: str type: str # 实体类型:Person, Organization, Technology, Product properties: Dict # 实体属性 description: str # 实体描述 @dataclass class Relation: """知识图谱关系""" source: str # 源实体名 target: str # 目标实体名 relation_type: str # 关系类型:belongs_to, uses, depends_on, manages properties: Dict class KnowledgeGraphManager: """知识图谱管理器""" def __init__(self, kg_client): self.client = kg_client async def add_entity(self, entity: Entity) -> bool: """添加实体到知识图谱""" query = """ MERGE (e:Entity {name: $name}) SET e.type = $type, e.description = $description, e += $properties RETURN e """ params = { "name": entity.name, "type": entity.type, "description": entity.description, "properties": entity.properties } result = await self.client.execute(query, params) return result is not None async def add_relation(self, relation: Relation) -> bool: """添加关系到知识图谱""" query = """ MATCH (s:Entity {name: $source}) MATCH (t:Entity {name: $target}) MERGE (s)-[r:RELATION {type: $rel_type}]->(t) SET r += $properties RETURN r """ params = { "source": relation.source, "target": relation.target, "rel_type": relation.relation_type, "properties": relation.properties } result = await self.client.execute(query, params) return result is not None async def get_entity(self, name: str) -> Optional[Entity]: """查询实体详情""" query = """ MATCH (e:Entity {name: $name}) RETURN e.name AS name, e.type AS type, e.description AS description, properties(e) AS properties """ result = await self.client.execute(query, {"name": name}) if result: return Entity(**result[0]) return None async def get_subgraph(self, center: str, max_hops: int = 2) -> Dict: """获取以某实体为中心的子图""" query = """ MATCH (center:Entity {name: $name}) CALL apoc.path.subgraphAll(center, {maxLevel: $hops}) YIELD nodes, relationships RETURN nodes, relationships """ result = await self.client.execute(query, {"name": center, "hops": max_hops}) return result async def delete_entity(self, name: str) -> bool: """删除实体及其所有关系""" query = """ MATCH (e:Entity {name: $name}) DETACH DELETE e """ result = await self.client.execute(query, {"name": name}) return True
使用依赖注入和配置管理,实现组件的灵活替换:
from dataclasses import dataclass from typing import Optional import json @dataclass class GraphRAGConfig: """GraphRAG 系统配置""" # 图数据库配置 neo4j_uri: str = "bolt://localhost:7687" neo4j_user: str = "neo4j" neo4j_password: str = "password" # 向量数据库配置 vector_db_type: str = "faiss" # faiss | milvus milvus_uri: Optional[str] = None embedding_model: str = "text-embedding-ada-002" embedding_dim: int = 1536 # 检索配置 graph_weight: float = 0.6 default_top_k: int = 10 max_hops: int = 3 # LLM 配置 llm_provider: str = "openai" # openai | local llm_model: str = "gpt-4" llm_temperature: float = 0.7 # 缓存配置 cache_enabled: bool = True cache_ttl: int = 3600 # 秒 redis_uri: Optional[str] = None @classmethod def from_file(cls, path: str) -> "GraphRAGConfig": """从 JSON 配置文件加载""" with open(path, 'r') as f: data = json.load(f) return cls(**data) class GraphRAGFactory: """工厂类:根据配置创建系统组件""" def __init__(self, config: GraphRAGConfig): self.config = config def create_kg_client(self) -> KnowledgeGraphClient: """创建知识图谱客户端""" if "neo4j" in self.config.neo4j_uri: from neo4j import AsyncGraphDatabase return Neo4jClient( uri=self.config.neo4j_uri, user=self.config.neo4j_user, password=self.config.neo4j_password ) raise ValueError(f"Unsupported graph database: {self.config.neo4j_uri}") def create_vector_store(self) -> VectorStoreClient: """创建向量存储客户端""" if self.config.vector_db_type == "faiss": return FAISSClient(dimension=self.config.embedding_dim) elif self.config.vector_db_type == "milvus": return MilvusClient(uri=self.config.milvus_uri) raise ValueError(f"Unsupported vector store: {self.config.vector_db_type}") def create_query_processor(self) -> QueryProcessor: """创建查询处理器""" return QueryProcessor( kg_client=self.create_kg_client(), vector_client=self.create_vector_store(), graph_weight=self.config.graph_weight )
GraphRAG 系统的容器化部署需要考虑各层的资源需求和通信方式:
# docker-compose.yml 示例 version: '3.8' services: # 数据接入服务(CPU/GPU 密集型) data-ingestion: build: ./services/data-ingestion deploy: resources: reservations: devices: - capabilities: [gpu] # NER 模型推理需要 GPU environment: - NEO4J_URI=bolt://neo4j:7687 - MILVUS_URI=milvus:19530 depends_on: - neo4j - milvus # 检索引擎服务(延迟敏感) retrieval-engine: build: ./services/retrieval-engine deploy: resources: limits: memory: 4G environment: - NEO4J_URI=bolt://neo4j:7687 - MILVUS_URI=milvus:19530 - REDIS_URI=redis:6379 ports: - "8001:8001" depends_on: - neo4j - milvus - redis # API 网关服务 api-gateway: build: ./services/api-gateway ports: - "8000:8000" environment: - RETRIEVAL_URL=http://retrieval-engine:8001 - LLM_API_KEY=${LLM_API_KEY} # 基础设施 neo4j: image: neo4j:5.15 ports: - "7474:7474" - "7687:7687" volumes: - neo4j_data:/data milvus: image: milvusdb/milvus:v2.3 ports: - "19530:19530" redis: image: redis:7-alpine ports: - "6379:6379" volumes: neo4j_data:
部署要点:
图数据库选型:Neo4j 是社区生态最成熟的图数据库,拥有丰富的查询语言 Cypher 和完善的可视化工具,适合中小规模的 GraphRAG 系统。如果预期图谱规模超过千万级节点,可以考虑 Neptune 或 NebulaGraph 等分布式图数据库。
向量索引方案:FAISS 在单机场景下性能优异,且与 Python 生态深度集成。对于需要分布式向量检索的场景,可以考虑 Milvus 或 Weaviate。向量索引的维度通常取 768(BERT-base)或 1536(text-embedding-ada-002),需要在精度和性能之间权衡。
异步处理架构:数据处理管道(文档解析→实体识别→关系抽取→图谱写入)天然适合异步处理。建议使用消息队列(如 RabbitMQ 或 Kafka)作为管道各阶段的缓冲,配合工作节点实现弹性伸缩。
本节从工程实现的角度,详细展示了 GraphRAG 系统的五层架构设计。通过查询处理引擎、知识图谱管理器、配置管理和容器化部署的完整代码示例,读者可以建立了一个可在生产环境中使用的系统骨架。
关键要点:
关键词:系统架构, 分层设计, 查询引擎, 知识图谱管理, 容器化部署, 依赖注入, GraphRAG工程
难度:进阶
预计阅读:30 分钟