📦 RAGFlow rag/ 模块源码分析报告

包含模块依赖关系图、主要功能描述和关键源码分析

生成时间: 2026-06-16 | 版本: ragflow-0.25.6

📊 概览统计

9
子包 (Subpackages)
90+
Python 模块
200+
核心类 & 函数
50+
LLM 提供商支持
15
文档解析器
3
GraphRAG 策略

1. 整体架构概览

rag/ 是 RAGFlow 的核心 RAG 引擎包,包含从文档解析、文本分块、向量化、 检索排序、到知识图谱构建的完整链路。其架构分为以下层次:

🏗️ 分层架构图

graph TB
    subgraph APP["应用层 (rag/app/)"]
        naive["naive.py
通用文档"] paper["paper.py
学术论文"] book["book.py
书籍"] qa["qa.py
问答对"] resume["resume.py
简历解析"] laws["laws.py
法律文档"] manual["manual.py
用户手册"] table["table.py
表格数据"] picture["picture.py
图像/视频"] presentation["presentation.py
幻灯片"] email_p["email.py
邮件"] tag["tag.py
标签对"] one["one.py
单块全文"] audio["audio.py
音频转录"] end subgraph FLOW["流水线层 (rag/flow/)"] File["File
文档加载"] Parser["Parser
文档解析"] Chunker["Chunker
文本分块"] Tokenizer["Tokenizer
分词与向量化"] Extractor["Extractor
LLM信息提取"] end subgraph NLP["NLP 层 (rag/nlp/)"] Dealer["Dealer
搜索检索"] FulltextQ["FulltextQueryer
全文查询构建"] Tokenizer2["RagTokenizer
分词器"] Synonym["Synonym
同义词"] TermWeight["TermWeight
词权重"] end subgraph LLM["LLM 抽象层 (rag/llm/)"] Chat["ChatModel
对话模型"] Embedding["EmbeddingModel
嵌入模型"] Rerank["RerankModel
重排序模型"] CV["CVModel
视觉模型"] TTS["TTSModel
语音合成"] OCR["OCRModel
OCR解析"] Seq2txt["Seq2txtModel
语音转文本"] end subgraph GRAPHRAG["GraphRAG 层 (rag/graphrag/)"] KGSearch["KGSearch
图谱检索"] EntRes["EntityResolution
实体消歧"] General["General
微软风格"] Light["Light
LightRAG风格"] NER["NER
spaCy提取"] end subgraph INFRA["基础设施层 (rag/utils/ + rag/svr/)"] Storage["存储连接器
ES/Infinity/Redis/MinIO/S3/..."] TaskExec["TaskExecutor
任务执行器"] SyncDS["SyncDataSource
数据源同步"] end APP --> FLOW FLOW --> NLP FLOW --> LLM NLP --> LLM NLP --> INFRA GRAPHRAG --> LLM GRAPHRAG --> NLP GRAPHRAG --> INFRA LLM --> INFRA

2. 模块依赖关系图

🔗 rag/ 包级依赖关系

核心模块 处理流水线 文档解析器 检索与图谱 基础设施
graph LR
    %% 核心模块
    rag_nlp["rag.nlp
搜索/分词/权重"]:::core rag_llm["rag.llm
LLM模型抽象"]:::core rag_prompts["rag.prompts
提示词管理"]:::core %% 流水线 rag_flow["rag.flow
文档处理流水线"]:::pipeline %% 解析器 rag_app["rag.app
文档类型解析器"]:::parser %% 高级检索 rag_graphrag["rag.graphrag
知识图谱RAG"]:::retrieval rag_advanced["rag.advanced_rag
深度研究"]:::retrieval rag_raptor["rag.raptor
RAPTOR摘要树"]:::retrieval %% 基础设施 rag_utils["rag.utils
工具/存储/缓存"]:::infra rag_svr["rag.svr
服务/任务/同步"]:::infra rag_benchmark["rag.benchmark
检索基准测试"]:::infra %% 依赖关系 rag_app --> rag_nlp rag_app --> rag_utils rag_flow --> rag_app rag_flow --> rag_nlp rag_flow --> rag_llm rag_flow --> rag_prompts rag_flow --> rag_utils rag_graphrag --> rag_nlp rag_graphrag --> rag_llm rag_graphrag --> rag_prompts rag_graphrag --> rag_utils rag_advanced --> rag_nlp rag_advanced --> rag_llm rag_advanced --> rag_prompts rag_advanced --> rag_utils rag_raptor --> rag_llm rag_raptor --> rag_graphrag rag_raptor --> rag_utils rag_svr --> rag_app rag_svr --> rag_flow rag_svr --> rag_llm rag_svr --> rag_graphrag rag_svr --> rag_raptor rag_svr --> rag_utils rag_benchmark --> rag_nlp rag_benchmark --> rag_llm classDef core fill:#dbeafe,stroke:#2563eb,color:#1e40af classDef pipeline fill:#dcfce7,stroke:#16a34a,color:#166534 classDef parser fill:#ffedd5,stroke:#ea580c,color:#9a3412 classDef retrieval fill:#f3e8ff,stroke:#9333ea,color:#6b21a8 classDef infra fill:#fee2e2,stroke:#dc2626,color:#991b1b

📄 文档处理流水线 (rag/flow/) 内部依赖

graph LR
    File["File
文档加载入口"] --> Parser["Parser
格式解析分发"] Parser --> TokenChunker["TokenChunker
Token/分隔符分块"] Parser --> TitleChunker["TitleChunker
标题层级分块"] TokenChunker --> Tokenizer2["Tokenizer
分词 & 向量嵌入"] TitleChunker --> Tokenizer2 Tokenizer2 --> Extractor2["Extractor
LLM元数据/TOC提取"] Parser --> DeepDOC["DeepDOC
PDF/Office解析"] Parser --> MinerU["MinerU OCR"] Parser --> Docling["Docling"] Parser --> TCADP["TCADP
腾讯云解析"] Parser --> PaddleOCR["PaddleOCR"] Parser --> Vision["Vision LLM
图像描述"] Parser --> Seq2txtModel["Seq2txt
语音转文本"] style File fill:#dbeafe,stroke:#2563eb style Parser fill:#dcfce7,stroke:#16a34a style TokenChunker fill:#ffedd5,stroke:#ea580c style TitleChunker fill:#ffedd5,stroke:#ea580c style Tokenizer2 fill:#f3e8ff,stroke:#9333ea style Extractor2 fill:#fee2e2,stroke:#dc2626

3. 核心模块 rag/ 顶层

🏷️ rag/__init__.py 包初始化

最小化的包初始化文件,仅包含版权声明和注释掉的 beartype 类型检查代码。不导出任何公共 API。

⚙️ rag/settings.py 配置占位

空模块,仅包含版权声明。原本可能用于 RAG 应用级配置,目前未被使用(实际配置在 common/settings.py 中)。

📊 rag/benchmark.py 基准测试

主要类:Benchmark(kb_id)

检索基准测试框架,支持 MS MARCO v1.1、TriviaQA、MIRACL 等标准 IR 数据集。

核心方法:

方法功能
ms_marco_index()索引 MS MARCO 数据集 (parquet)
trivia_qa_index()索引 TriviaQA 数据集
miracl_index()索引 MIRACL 多语言数据集
_get_retrieval()对所有查询执行检索,返回 run 字典
save_results()保存 NDCG@10、QRELS、RUN 文件

关键依赖:common.*, rag.nlp.search, ranx (评估指标), pandas

🌳 rag/raptor.py RAPTOR 算法

主要类:RecursiveAbstractiveProcessing4TreeOrganizedRetrieval

实现 RAPTOR 递归摘要树算法,支持两种构建策略:

  • Classic RAPTOR: UMAP 降维 + GMM/AHC 聚类 → 递归生成摘要层
  • Psi Tree: 直接基于成对余弦相似度排名构建合并树,更快且无需降维

核心方法:

方法功能
_get_optimal_clusters()通过 BIC 选择最优 GMM 聚类数
_get_clusters_ahc()Ward 层次聚类 + 树状图间隙启发式
_build_psi_structure()Psi 合并树构建
_summarize_texts()LLM 摘要生成(含缓存)

关键依赖:umap, sklearn, rag.graphrag.utils (LLM/Embedding 缓存), rag.utils.raptor_utils

4. LLM 模型抽象层 rag/llm/

🏭 LLM 工厂模式与模型继承关系

graph TB
    Init["__init__.py
自动发现 & 注册
SupportedLiteLLMProvider 枚举"] subgraph Chat["chat_model.py"] ChatBase["Base(ABC)
OpenAI兼容基类"] LiteLLMBaseChat["LiteLLMBase(ABC)
LiteLLM路由基类"] ConcreteChat["30+ 具体实现
OpenAI, Anthropic, Bedrock
Ollama, vLLM, ..."] end subgraph Embedding["embedding_model.py"] EmbedBase["Base(ABC)
嵌入模型基类"] ConcreteEmb["30+ 具体实现
OpenAI, Jina, Cohere
BGE, QWen, ..."] end subgraph Rerank["rerank_model.py"] RerankBase["Base(ABC)
重排序基类"] ConcreteRerank["15+ 具体实现
Jina, Cohere, QWen
Voyage, ..."] end subgraph CV["cv_model.py"] CVBase["Base(ABC)
视觉模型基类"] GptV4["GptV4
OpenAI视觉"] ConcreteCV["20+ 具体实现
Gemini, Anthropic, Ollama
QWen(视频), ..."] end subgraph Other["其他模型"] OCRModel["ocr_model.py
MinerU/PaddleOCR/OpenDataLoader"] TTSModel["tts_model.py
10+ TTS实现"] Seq2txtModel["sequence2txt_model.py
15+ 语音转文本"] ToolDeco["tool_decorator.py
@tool 装饰器"] end Init --> Chat Init --> Embedding Init --> Rerank Init --> CV Init --> Other ChatBase --> ConcreteChat LiteLLMBaseChat --> ConcreteChat

🔌 llm/__init__.py — 自动发现与注册 工厂模式

通过 importlib + inspect 自动扫描子模块,提取 _FACTORY_NAME 并注册到工厂字典。

# 核心注册逻辑
MODULE_MAPPING = {
    "chat_model": ChatModel,
    "embedding_model": EmbeddingModel,
    "rerank_model": RerankModel,
    "tts_model": TTSModel,
    "cv_model": CvModel,
    "ocr_model": OcrModel,
    "sequence2txt_model": Seq2txtModel,
}

# 运行时自动发现所有带 _FACTORY_NAME 的类
for mod_name, factory_dict in MODULE_MAPPING.items():
    module = importlib.import_module(f"rag.llm.{mod_name}")
    for name, cls in inspect.getmembers(module, inspect.isclass):
        if hasattr(cls, '_FACTORY_NAME'):
            for n in cls._FACTORY_NAME if isinstance(...) else [cls._FACTORY_NAME]:
                factory_dict[n] = cls  # 注册: ChatModel["OpenAI"] = OpenAI_APIChat

💬 llm/chat_model.py — 对话模型 30+ 提供商

核心类层级:

  • Base(ABC) — OpenAI 兼容基类,提供流式/非流式对话、多轮工具调用、错误分类重试、reasoning 内容处理
  • LiteLLMBase(ABC) — LiteLLM 路由基类,通过 litellm.acompletion 统一接入 100+ 提供商
  • 30+ 具体实现(XinferenceChat, HuggingFaceChat, GoogleChat, BaiduYiyanChat, SparkChat...)

关键机制:

  • 错误分类:LLMErrorCode 枚举(RATE_LIMIT、AUTH_ERROR、SERVER_ERROR 等),自动重试
  • 模型策略:_apply_model_family_policies() 对 Qwen3/Kimi K2.5/GPT-5 等模型应用参数约束
  • 工具调用:async_chat_with_tools() 多轮自动重试工具调用

📐 llm/embedding_model.py — 嵌入模型 30+ 提供商

统一接口:encode(texts) -> (np.ndarray, token_count)encode_queries(text) -> (np.ndarray, token_count)

关键实现:

  • JinaMultiVecEmbed — 支持 v4 多向量(mean-pooling)和 v2/v3 单向量
  • PerplexityEmbed — 支持 contextualized embeddings (base64 int8 解码)
  • SILICONFLOWEmbed — BGE 模型自动截断至 256 tokens
  • BuiltinEmbed — 线程安全单例,自动加载 TEI 或 HuggingFace 模型

🎯 llm/rerank_model.py — 重排序模型 15+ 提供商

统一接口:similarity(query, texts) -> (np.ndarray, token_count)

所有实现均提供 _normalize_rank() 静态方法进行 min-max 分数归一化。

👁️ llm/cv_model.py — 视觉模型 20+ 提供商

统一接口:describe(image) -> (text, token_count)

亮点:

  • QWenCV — 额外支持视频处理(DashScope MultiModalConversation)
  • GeminiCV — 视频处理(20MB 阈值,Files API 处理大文件)
  • AnthropicCV — 支持 thinking_delta 流式推理
  • 所有基于 GptV4 的实现共享 _form_history 多模态消息构建逻辑

🛠️ llm/tool_decorator.py — 工具装饰器 轻量级

@tool
def search_web(query: str, max_results: int = 5) -> list[dict]:
    """Search the web for information.
    :param query: The search query string.
    :param max_results: Maximum number of results to return.
    """
    ...

# 自动生成 OpenAI function schema:
# {"type": "function", "function": {"name": "search_web",
#  "description": "Search the web...", "parameters": {...}}}

FunctionToolSession 将装饰函数桥接到对话模型的工具调用循环。

5. NLP 与搜索引擎 rag/nlp/

🔍 搜索引擎组件交互

graph TB
    User["用户查询"] --> Dealer["Dealer
search.py"] Dealer --> FulltextQ["FulltextQueryer
query.py
全文查询构建"] Dealer --> DocStore["DocStore
ES/Infinity/OB"] FulltextQ --> Tokenizer3["RagTokenizer
rag_tokenizer.py"] FulltextQ --> TermW["TermWeight
term_weight.py
词权重 + NER"] FulltextQ --> Syn["Synonym
synonym.py
同义词扩展"] FulltextQ --> Surname["Surname
surname.py
中文姓氏库"] Dealer --> EmbedModel["EmbeddingModel
向量编码"] Dealer --> RerankModel["RerankModel
精排重排序"] style Dealer fill:#dbeafe,stroke:#2563eb style FulltextQ fill:#dcfce7,stroke:#16a34a style DocStore fill:#fee2e2,stroke:#dc2626

🔎 nlp/search.py — 核心检索引擎 核心

主要类:Dealer(dataStore) — 多阶段检索引擎

检索流程:

  1. 初检:search() — 文本/向量/混合搜索 + 融合表达式
  2. 过滤:_prune_deleted_chunks() — 移除已删除文档的 chunks
  3. 重排序:rerank_by_model() — 外部 Cross-Encoder 精排
  4. 排序:按相似度阈值过滤,分页
  5. 聚合:group_docs — 按文档分组结果

核心方法一览:

方法功能
search()多模态初检(文本/向量/混合)
rerank_by_model()Cross-Encoder 模型精排
rerank()本地混合相似度重排
retrieval()完整检索管道入口
insert_citations()答案句级引用标注
retrieval_by_toc()TOC 感知检索
hybrid_similarity()向量+Token 混合相似度

📝 nlp/query.py — 全文查询构建 查询理解

主要类:FulltextQueryer(QueryBase)

将自然语言 query 转换为结构化全文查询:

  • 英文路径:词权重计算 → 同义词扩展 → Bigram 生成 → 字段加权
  • 中文路径:细粒度分词 → 同义词扩展 → 子词 n-gram → 位置加权

提供 hybrid_similarity()token_similarity() 用于检索后重排序。

✂️ nlp/rag_tokenizer.py — 分词器 Infinity 适配

主要类:RagTokenizer (继承 Infinity 的 C++ 分词器)

当使用 Infinity 引擎时,tokenize()fine_grained_tokenize() 为透传(Infinity 自行处理分词)。

模块级导出:tokenize, fine_grained_tokenize, tag, freq, tradi2simp

🔤 nlp/term_weight.py — 词权重 排序信号

类:Dealer — 组合多种信号计算词权重:NER 类型得分、词性标签权重、语料内词频、文档频率 IDF。

权重归一化至总和 1.0,用于查询构建时的字段加权。

6. 文档处理流水线 rag/flow/

⚙️ 流水线组件异步执行模型

sequenceDiagram
    participant Pipeline
    participant File
    participant Parser
    participant Chunker
    participant Tokenizer
    participant Extractor
    participant Redis

    Pipeline->>File: _invoke()
    File-->>Pipeline: {name, file}
    Pipeline->>Redis: callback(progress)
    Pipeline->>Parser: _invoke(name, file)
    Parser->>Parser: 检测文件类型
    Parser->>Parser: 分发到专用解析器
    Parser-->>Pipeline: {chunks/json/markdown/...}
    Pipeline->>Redis: callback(progress)
    Pipeline->>Chunker: _invoke(chunks)
    Chunker->>Chunker: TokenChunker/TitleChunker
    Chunker-->>Pipeline: {chunks}
    Pipeline->>Redis: callback(progress)
    Pipeline->>Tokenizer: _invoke(chunks)
    Tokenizer->>Tokenizer: 全文分词 + 向量嵌入
    Tokenizer-->>Pipeline: {chunks}
    Pipeline->>Redis: callback(progress)
    Pipeline->>Extractor: _invoke(chunks)
    Extractor->>Extractor: LLM TOC/元数据提取
    Extractor-->>Pipeline: {chunks}
    Pipeline->>Redis: callback(done)
    

🏭 flow/__init__.py — 自动组件发现 反射机制

通过 pkgutil.walk_packages 遍历所有子包,自动导入并注册所有公共类到 rag.flow 命名空间。

🧱 flow/base.py — 流水线基类 抽象基类

核心类:

  • ProcessParamBase(ComponentParamBase) — 参数基类,默认 timeout=100000000
  • ProcessBase(ComponentBase) — 组件基类,提供 async _invoke(**kwargs) 抽象方法

执行模型:超时保护 + 异常捕获 + 进度回调 + 耗时统计

@timeout(int(os.environ.get("COMPONENT_EXEC_TIMEOUT", 60*60)))
async def _invoke(self, **kwargs):
    ...

🚀 flow/pipeline.py — 流水线编排器 核心

主要类:Pipeline(Graph) — 继承 Agent 框架的 Graph 类

从 DSL JSON 解析组件图,沿 DAG 路径串行执行组件。每个组件通过 invoke() 异步调用,进度通过 Redis 日志记录。

执行路径:File → Parser → Chunker → Tokenizer → Extractor(线性 DAG)

取消机制:通过 Redis 检查取消标记,抛出 TaskCanceledException

📄 flow/parser/parser.py — 文档解析分发器 最大模块

类:Parser(ProcessBase) — 支持 12+ 文件格式

方法处理格式后端
_pdf()PDFDeepDOC/PlainText/MinerU/Docling/OpenDataLoader/TCADP/PaddleOCR/Vision
_docx()DOCXDeepDOC + 大纲/页眉页脚/TOC 去除
_doc()DOCApache Tika
_spreadsheet()XLSX/CSVTCADP/DeepDOC ExcelParser
_slides()PPT/PPTXTCADP/DeepDOC PPT Parser
_markdown()MD/MDX内建解析(含图片/表格)
_image()图片OCR + VLM 描述
_audio()音频Speech-to-Text 模型
_video()视频Image-to-Text 模型
_email()EML/MSG邮件头/正文/附件提取
_html()HTML含页眉/页脚去除
_epub()EPUB章节分解

✂️ flow/chunker/ — 文本分块 两种策略

TokenChunker:按 token 数或自定义分隔符分块,支持重叠。对表格/图片自动附加上下文窗口。

TitleChunker:基于标题层级的分块:

  • HierarchyTitleChunker — 树形层级分块(DFS 路径提取)
  • GroupTitleChunker — 扁平分组分块(Token 约束合并)
  • BaseTitleChunker — 共享基类,标题级别识别(PDF 书签/正则匹配)

🔢 flow/tokenizer/tokenizer.py — 分词与向量化

对每个 chunk 执行全文分词(标题/内容/问题关键词/重要关键词)并生成向量嵌入。支持文件名加权向量计算和批量处理速率限制。

🤖 flow/extractor/extractor.py — LLM 信息提取

Extractor(ProcessBase, LLM) — 双重继承自流水线基类和 Agent LLM 组件。

两种模式:

  • TOC 生成:field_name == "toc" 时,调用 run_toc_from_text() 生成目录
  • 元数据提取:自定义 LLM Prompt 提取指定字段

7. 文档类型解析器 rag/app/

📑 文档解析器继承与分发

graph TB
    subgraph "rag.app 解析器"
        naive["naive.py
通用解析
chunk() 分发器"] paper["paper.py
学术论文"] book["book.py
书籍"] qa2["qa.py
问答对"] resume["resume.py
简历"] laws["laws.py
法律"] manual["manual.py
手册"] table2["table.py
表格"] picture["picture.py
图片/视频"] presentation["presentation.py
幻灯片"] email2["email.py
邮件"] tag2["tag.py
标签对"] one2["one.py
全文不分割"] audio["audio.py
音频"] end subgraph "DeepDoc 基类" PdfParser["PdfParser"] DocxParser["DocxParser"] ExcelParser["ExcelParser"] HtmlParser["HtmlParser"] TxtParser["TxtParser"] PptParser["RAGFlowPptParser"] end naive --> PdfParser naive --> DocxParser naive --> ExcelParser naive --> HtmlParser paper --> PdfParser book --> PdfParser book --> DocxParser qa2 --> PdfParser qa2 --> ExcelParser resume --> PdfParser laws --> PdfParser laws --> DocxParser manual --> PdfParser manual --> DocxParser table2 --> ExcelParser presentation --> PdfParser presentation --> PptParser one2 --> PdfParser one2 --> ExcelParser style naive fill:#dbeafe,stroke:#2563eb

各解析器对比:

模块文档类型分块策略特色
naive.py通用 (15+ 格式)naive_merge token 合并最广泛格式支持,主分发器
paper.py学术论文 (PDF)标题层级分块提取标题/作者/摘要,双栏布局处理
book.py书籍 (5 格式)hierarchical_mergeTOC 去除,章节层级分块
qa.py问答对 (6 格式)每个 Q&A 对为一块问题-答案检测(PDF 项目符号/MD 标题/DOCX 样式)
resume.py简历 (3 格式)LLM 结构化提取双路径融合 + 并行 LLM + 4 阶段后处理(arXiv:2510.09722)
laws.py法律 (6 格式)tree_merge层级树形合并
manual.py手册 (2 格式)节 ID 合并PDF 书签标题级别检测
table.py表格 (3 格式)每行一块多级表头/合并单元格/列类型推断/VLM 图片描述
picture.py图片/视频每文件一块OCR + VLM 双重描述
presentation.py幻灯片 (2 格式)每页一块页面缩略图,位置排序
email.py邮件 (EML)头+正文合并附件递归解析
tag.py标签对 (3 格式)每行一块标签关键词字段
one.py任意 (6 格式)全文一块最简模式
audio.py音频 (12+ 格式)每文件一块Whisper 等语音转文本

8. GraphRAG 知识图谱 rag/graphrag/

🕸️ GraphRAG 流水线与检索

graph TB
    subgraph Index["图谱构建 (graphrag/general/index.py)"]
        Docs["文档集合"] --> SubGraph["generate_subgraph()
逐文档构建子图"] SubGraph --> Merge["merge_subgraph()
合并到全局图"] Merge --> EntRes2["resolve_entities()
LLM实体消歧"] EntRes2 --> Community["extract_community()
Leiden社区发现 + LLM报告"] end subgraph Extractors["实体/关系提取器"] GeneralExt["General
graph_extractor.py
微软GraphRAG风格
带gleaning迭代"] LightExt["Light
graph_extractor.py
LightRAG风格
简化prompt"] NerExt["NER
graph_extractor.py
spaCy提取
无需LLM"] MindMapExt["MindMap
mind_map_extractor.py
思维导图生成"] end subgraph Retrieval["图谱检索 (graphrag/search.py)"] KGSearch2["KGSearch(Dealer)"] Rewrite["query_rewrite()
查询改写为类型词+实体词"] EntRet["get_relevant_ents_by_keywords()
实体检索"] RelRet["get_relevant_relations_by_txt()
关系检索"] NHop["n-hop路径评分"] ComRet["_community_retrieval_()
社区报告检索"] end Docs --> Extractors Extractors --> SubGraph KGSearch2 --> Rewrite Rewrite --> EntRet Rewrite --> RelRet EntRet --> NHop RelRet --> NHop EntRet --> ComRet NHop --> KGSearch2 ComRet --> KGSearch2 style Index fill:#dcfce7,stroke:#16a34a style Extractors fill:#f3e8ff,stroke:#9333ea style Retrieval fill:#dbeafe,stroke:#2563eb

🏭 graphrag/general/index.py — 构建编排器 核心

主函数:run_graphrag_for_kb() — 完整的 GraphRAG 构建流水线

  1. 子图生成:generate_subgraph() 逐文档并行提取实体/关系
  2. 子图合并:merge_subgraph() 合并到全局 NetworkX 图
  3. 实体消歧:resolve_entities() LLM 判断同名实体是否指代同一事物
  4. 社区报告:extract_community() Leiden 聚类 + LLM 生成社区摘要

可恢复性:通过 Redis 阶段标记 (phase_markers.py) 支持中断后恢复。

三种提取策略:

  • general — 微软 GraphRAG 风格,LLM gleaning 迭代提取
  • light — LightRAG 风格,简化 Prompt
  • ner — spaCy NER + 共现关系,无需 LLM

🔍 graphrag/search.py — 图谱检索 多源融合

类:KGSearch(Dealer) — 继承自搜索引擎,扩展图谱检索能力

检索流程:查询改写 → 实体关键词检索 → 关系检索 → n-hop 路径扩展 → 社区报告检索

使用 xxhash 进行 LLM 结果缓存,避免重复调用。

🧬 graphrag/general/leiden.py — Leiden 聚类 社区发现

实现层次 Leiden 聚类算法,用于图谱社区发现。使用 graspologic 库生成加权层次社区结构,支持 stable_largest_connected_component 和确定性图排序。

🧠 graphrag/ner/graph_extractor.py — NER 提取 无 LLM

基于 spaCy 的实体/关系提取,无需调用 LLM:

  • MGranRAG 关键词堆叠:3 轮叠加(连字符/撇号合并 → 大写词合并 → 名词/数字合并)
  • LinearRAG 共现关系:同句共现或依存路径连接

9. 提示词管理 rag/prompts/

📋 prompts/template.py — 模板加载器

load_prompt(name: str) -> str — 从 prompts/ 目录加载 .md 模板文件,带内存缓存。

🤖 prompts/generator.py — 提示词生成中心 30+ 函数

所有 LLM 交互的提示词构建集中在此模块,使用 Jinja2 沙盒模板引擎。

关键功能分组:

功能组函数用途
检索增强kb_prompt(), citation_prompt(), citation_plus()知识库片段格式化,引用标注
查询理解keyword_extraction(), question_proposal(), full_question()关键词提取,问题生成,问题改写
多语言cross_languages()查询多语言翻译
内容标签content_tagging()自动内容标签
Agent 工具analyze_task_async(), next_step_async(), reflect_async()ReAct Agent 循环
记忆管理tool_call_summary(), rank_memories_async()工具调用摘要,记忆排序
TOC 生成detect_table_of_contents(), extract_table_of_contents(), run_toc_from_text()目录检测/提取/生成
元数据gen_metadata(), gen_meta_filter()元数据生成,元数据过滤
深度研究sufficiency_check(), multi_queries_gen()信息充分性检查,多查询生成

技术细节:gen_json() 使用 json_repair 修复 LLM 输出的畸形 JSON,支持重试。所有异步函数使用 chat_mdl.async_chat() 调用 LLM。

10. 工具与基础设施 rag/utils/

🗄️ 存储连接器层级

graph TB
    subgraph "文档存储 (DocStore)"
        ES["es_conn.py
ESConnection
(单例)"] Infinity["infinity_conn.py
InfinityConnection
(单例)"] OS["opensearch_conn.py
OSConnection
(单例)"] OB["ob_conn.py
OBConnection
(单例) OceanBase"] end subgraph "对象存储 (Object Store)" MinIO["minio_conn.py
RAGFlowMinio
(单例)"] S3["s3_conn.py
RAGFlowS3
(单例)"] OSS["oss_conn.py
RAGFlowOSS
(单例) 阿里云"] AzureSAS["azure_sas_conn.py
Azure SAS"] AzureSPN["azure_spn_conn.py
Azure Data Lake"] GCS["gcs_conn.py
Google Cloud"] OpenDAL["opendal_conn.py
OpenDAL 通用"] end subgraph "缓存与消息" Redis["redis_conn.py
RedisDB + RedisDistributedLock
(单例)"] Encrypted["encrypted_storage.py
透明加密包装"] end subgraph "外部集成" Tavily2["tavily_conn.py
Tavily搜索"] TTSCache["tts_cache.py
TTS缓存"] end subgraph "工具函数" FileUtils["file_utils.py
文件提取/链接/HTML"] RaptorUtils["raptor_utils.py
RAPTOR配置与标记"] ImageUtils["base64_image.py
图片存储"] LazyImage["lazy_image.py
延迟加载"] TableMeta["table_es_metadata.py
表格元数据聚合"] end style ES fill:#fee2e2,stroke:#dc2626 style Infinity fill:#fee2e2,stroke:#dc2626 style Redis fill:#dbeafe,stroke:#2563eb style MinIO fill:#dcfce7,stroke:#16a34a

🔌 文档存储连接器

模块后端关键特性
es_conn.pyElasticsearch全文/向量/融合搜索,Painless 脚本 Pagerank 调整,聚合
infinity_conn.pyInfinity多表搜索,字段名映射,DataFrame 字段提取
opensearch_conn.pyOpenSearchKNN DSL 搜索,布尔查询
ob_conn.pyOceanBase全文+向量混合搜索(融合退化为独立),完整 SQL 模式定义

🗄️ utils/redis_conn.py — Redis 工具 核心基础设施

类:RedisDB (单例) — 完整 Redis 操作封装

功能:KV 操作、集合、有序集合、自增 ID、Lua 脚本事务、Stream 消息队列(生产者/消费者)、未确认消息重处理。

类:RedisDistributedLock — 分布式锁(delete-if-equal 安全释放)

11. 高级 RAG rag/advanced_rag/

🌲 tree_structured_query_decomposition_retrieval.py DeepResearcher

类:TreeStructuredQueryDecompositionRetrieval (别名 DeepResearcher)

实现树形查询分解深度研究:

  1. 检索:从 KB + Web (Tavily) + KG 三源检索
  2. 充分性检查:LLM 判断信息是否足以回答问题
  3. 递归分解:若不充分,将问题分解为子问题,递归深度研究(有深度限制)
  4. 去重合并:线程安全的 chunk 信息合并

关键依赖:rag.prompts.generator (充分性检查/多查询生成), rag.utils.tavily_conn, rag.nlp.search

12. 服务运行器 rag/svr/

svr/task_executor.py — 任务执行器 核心

Redis Stream 消费者,处理所有文档摄入任务:

# 并发控制信号量
task_limiter       # 任务级并发限制
chunk_limiter      # 分块并发限制
embed_limiter      # 嵌入生成速率限制
minio_limiter      # MinIO 并发限制
kg_limiter         # 知识图谱并发限制

任务类型分发:

任务类型处理函数描述
MEMORYbuild_chunks()标准文档分块+嵌入
DATAFLOWrun_dataflow()流水线执行
RAPTORrun_raptor_for_kb()RAPTOR 摘要树
GRAPHRAGrun_graphrag_for_kb()知识图谱构建
MINDMAPrun_mindmap_for_kb()思维导图生成

RAPTOR 任务特殊处理:支持检查点(checkpointing)、清理已完成摘要 chunk、取消时的回滚删除。

🔄 svr/sync_data_source.py — 数据源同步 25+ 连接器

类层级:SyncBase_BlobLikeBase → 具体连接器

支持的数据源:S3, R2, OCI Storage, Google Cloud Storage, RSS, Confluence, Notion, Discord, Gmail, Dropbox, GoogleDrive, Jira, SharePoint, Slack, Teams, WebDAV, Moodle, BOX, Airtable, Asana, Github, Gitlab, Bitbucket, SeaFile, DingTalk, IMAP, Zendesk, MySQL, PostgreSQL, REST API。

同步机制:基于指纹的变更检测 + 定时轮询调度。

13. 数据处理全流程

🔄 端到端文档处理数据流

flowchart TB
    A["📥 文档上传"] --> B{"文档类型检测"}
    B -->|"PDF/DOCX/..."| C["rag.app.*
类型专用解析器"] C --> D["rag.flow.parser
统一解析"] D --> E["rag.flow.chunker
TokenChunker / TitleChunker"] E --> F["rag.flow.tokenizer
分词 + 向量嵌入"] F --> G["📊 向量数据库
ES / Infinity / OceanBase"] E --> H["rag.raptor
RAPTOR摘要树
(可选)"] H --> G C --> I["rag.graphrag
知识图谱构建
(可选)"] I --> J["🕸️ 图谱存储
实体 + 关系 + 社区报告"] J --> G K["🔍 用户查询"] --> L["rag.nlp.search.Dealer
多阶段检索"] L --> M["rag.nlp.query
全文查询构建"] L --> N["rag.llm.Embedding
向量编码"] L --> O["rag.llm.Rerank
重排序"] L --> P["rag.graphrag.KGSearch
图谱检索"] L --> Q["rag.advanced_rag
深度研究"] M --> G N --> G P --> J L --> R["📋 检索结果
+ 引用标注"] style A fill:#dbeafe,stroke:#2563eb style G fill:#dcfce7,stroke:#16a34a style J fill:#f3e8ff,stroke:#9333ea style R fill:#fee2e2,stroke:#dc2626

关键设计模式总结

模式应用场景实现方式
工厂模式LLM 模型实例化rag.llm.__init__ 自动发现 _FACTORY_NAME 并注册
策略模式文档解析/分块/GraphRAG 提取根据配置选择不同策略类(General/Light/NER)
管道模式文档处理流水线rag.flow.Pipeline DAG 组件串行执行
单例模式存储连接器@singleton 装饰器确保全局唯一连接
装饰器模式工具函数注册@tool 自动生成 OpenAI function schema
缓存模式LLM/Embedding 结果Redis xxhash 缓存(get_llm_cache / set_llm_cache
分布式锁并发任务协调RedisDistributedLock + 阶段标记

核心数据流路径

场景路径
标准文档检索文档 → app/* → flow/parser → flow/chunker → flow/tokenizer → ES/Infinity → nlp/search → 结果
RAPTOR 分层检索文档 → app/* → flow/* → raptor.py (递归摘要) → ES/Infinity → nlp/search (含摘要层)
GraphRAG 图谱检索文档 → app/* → graphrag/general/index.py (实体/关系/社区) → graphrag/search.py → 结果
深度研究查询 → advanced_rag (分解+检索+检查) → KB/Web/KG 多源 → 合成答案
数据源同步外部数据源 → svr/sync_data_source.py → 文档摄入管道