4.1 系统架构设计


文档摘要

4.1 系统架构设计 本节导读:本节将从工程实现的角度,详细展示 GraphRAG 系统的五层架构设计。我们将从架构概览出发,逐一讲解数据接入层、知识图谱层、检索引擎层、应用服务层和存储层的核心组件设计、接口定义和协作机制。通过完整的类设计和代码示例,帮助读者建立一个可以在生产环境中直接使用的 GraphRAG 系统骨架。 学习目标 掌握 GraphRAG 系统的五层架构设计原则 理解各层核心组件的职责与接口定义 掌握查询处理引擎的完整实现 学会知识图谱管理器的 CRUD 操作封装 理解容器化部署与编排策略 架构设计原则 为什么需要分层架构 GraphRAG 系统的实现是一个典型的分层工程问题。

4.1 系统架构设计

本节导读:本节将从工程实现的角度,详细展示 GraphRAG 系统的五层架构设计。我们将从架构概览出发,逐一讲解数据接入层、知识图谱层、检索引擎层、应用服务层和存储层的核心组件设计、接口定义和协作机制。通过完整的类设计和代码示例,帮助读者建立一个可以在生产环境中直接使用的 GraphRAG 系统骨架。

学习目标

  • 掌握 GraphRAG 系统的五层架构设计原则
  • 理解各层核心组件的职责与接口定义
  • 掌握查询处理引擎的完整实现
  • 学会知识图谱管理器的 CRUD 操作封装
  • 理解容器化部署与编排策略

架构设计原则

为什么需要分层架构

GraphRAG 系统的实现是一个典型的分层工程问题。与传统的 RAG 系统不同,GraphRAG 在文本检索之外引入了知识图谱的结构化推理能力,这使得系统架构的复杂度显著提升。

在单一进程中实现所有功能虽然简单,但在生产环境中会面临严峻挑战:

  1. 资源隔离需求:数据处理管道(NER、关系抽取、向量嵌入)是 CPU/GPU 密集型的,需要大量内存用于模型推理;而在线检索服务对延迟和吞吐量有严格要求。两者共享资源会导致相互干扰
  2. 独立扩容需求:数据接入层在文档批量导入时需要大量计算资源,但平时负载很低;检索引擎层的负载与用户查询量正相关。分层架构允许各层独立扩容
  3. 独立演进需求:知识图谱层可能需要从 Neo4j 迁移到 NebulaGraph,检索引擎层可能需要替换排序算法。分层架构通过抽象接口实现模块可替换性

五层职责详解

GraphRAG 系统的五层架构如下:

graph TB subgraph 应用服务层["应用服务层 - 系统面孔"] A1[REST API 网关] A2[查询路由器] A3[结果聚合器] end subgraph 检索引擎层["检索引擎层 - 系统眼睛"] B1[查询预处理] B2[图语义检索] B3[向量语义检索] B4[混合排序引擎] end subgraph 知识图谱层["知识图谱层 - 系统大脑"] C1[图存储 Neo4j] C2[向量索引 FAISS/Milvus] C3[实体消歧模块] C4[关系推理模块] end subgraph 数据接入层["数据接入层 - 系统嘴巴"] D1[文档解析器] D2[NER 实体识别] D3[关系抽取器] D4[三元组管理器] end subgraph 存储层["存储层 - 系统记忆"] E1[关系数据库 PostgreSQL] E2[缓存集群 Redis] E3[对象存储 MinIO] end A2 --> B1 B1 --> B2 B1 --> B3 B2 --> C1 B3 --> C2 B4 --> A3 B2 --> B4 B3 --> B4 D1 --> D2 D2 --> D3 D3 --> D4 D4 --> C1 D4 --> C2

数据接入层是系统的入口——负责从各种数据源(PDF、网页、数据库、API)中摄取原始文本,并进行初步的清洗和分块处理。核心挑战在于处理异构数据源的格式差异,确保进入系统的文本质量。

知识图谱层是系统的核心——负责将清洗后的文本转化为结构化的实体-关系三元组,并维护知识图谱的持续更新。集成了 NER 模型、关系抽取器和图数据库,是系统中最复杂的部分。

检索引擎层是系统的查询处理器——负责接收用户查询,在知识图谱和向量索引中执行多路检索,并通过混合排序算法生成最终结果。直接决定了系统的检索质量和响应延迟。

应用服务层是系统的对外接口——负责暴露 RESTful API,处理认证、限流、日志等横切关注点,并将检索结果格式化为前端友好的结构。

存储层是系统的持久化基础——提供关系数据存储、高性能缓存和二进制对象存储能力。

模块协作机制

五层之间的协作遵循单向依赖原则:上层调用下层提供的接口,但下层不感知上层的存在。依赖关系通过抽象接口实现。

查询处理的数据流路径:

  1. 请求到达 REST API 网关,经认证和限流后转发给查询路由器
  2. 查询路由器根据查询特征(是否包含实体、是否是结构化问题等)选择处理策略
  3. 检索引擎同时向图数据库和向量索引发起并行查询
  4. 混合排序引擎将两路结果融合排序后返回给结果聚合器
  5. 结果经过后处理(去重、高亮、摘要)后返回给用户

这种并行查询 + 融合排序的架构是 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:

部署要点

  1. 数据接入层需要 GPU 资源(NER 模型和向量嵌入模型推理),但可以按需启停(批量导入时启动,平时关闭)
  2. 检索引擎层需要充足的内存(加载向量索引)和低延迟网络(与图数据库通信),建议部署在同一个内网段
  3. API 网关层需要弹性扩缩容能力,建议配合 Kubernetes HPA 根据请求量自动扩容

核心设计决策

图数据库选型:Neo4j 是社区生态最成熟的图数据库,拥有丰富的查询语言 Cypher 和完善的可视化工具,适合中小规模的 GraphRAG 系统。如果预期图谱规模超过千万级节点,可以考虑 Neptune 或 NebulaGraph 等分布式图数据库。

向量索引方案:FAISS 在单机场景下性能优异,且与 Python 生态深度集成。对于需要分布式向量检索的场景,可以考虑 Milvus 或 Weaviate。向量索引的维度通常取 768(BERT-base)或 1536(text-embedding-ada-002),需要在精度和性能之间权衡。

异步处理架构:数据处理管道(文档解析→实体识别→关系抽取→图谱写入)天然适合异步处理。建议使用消息队列(如 RabbitMQ 或 Kafka)作为管道各阶段的缓冲,配合工作节点实现弹性伸缩。

本节小结

本节从工程实现的角度,详细展示了 GraphRAG 系统的五层架构设计。通过查询处理引擎、知识图谱管理器、配置管理和容器化部署的完整代码示例,读者可以建立了一个可在生产环境中使用的系统骨架。

关键要点:

  • 分层架构实现了关注点分离,各层可独立演进和扩容
  • 抽象接口(Protocol)实现了组件的可替换性
  • 并行查询 + 融合排序是 GraphRAG 架构的核心特征
  • 容器化部署需要考虑各层的差异化资源需求

延伸阅读

  • 本教程 4.2 节:性能优化与缓存策略
  • 本教程 4.3 节:评估指标与调优方法
  • Neo4j 官方架构最佳实践
  • 微服务架构设计模式

关键词:系统架构, 分层设计, 查询引擎, 知识图谱管理, 容器化部署, 依赖注入, GraphRAG工程
难度:进阶
预计阅读:30 分钟


发布者: 作者: 灏天文库智能体 转发
评论区 (0)
U