rag/ 模块源码分析报告包含模块依赖关系图、主要功能描述和关键源码分析
生成时间: 2026-06-16 | 版本: ragflow-0.25.6
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
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
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
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 递归摘要树算法,支持两种构建策略:
核心方法:
| 方法 | 功能 |
|---|---|
_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
rag/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+ 提供商关键机制:
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 tokensBuiltinEmbed — 线程安全单例,自动加载 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 将装饰函数桥接到对话模型的工具调用循环。
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) — 多阶段检索引擎
检索流程:
search() — 文本/向量/混合搜索 + 融合表达式_prune_deleted_chunks() — 移除已删除文档的 chunksrerank_by_model() — 外部 Cross-Encoder 精排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 转换为结构化全文查询:
提供 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,用于查询构建时的字段加权。
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=100000000ProcessBase(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() | DeepDOC/PlainText/MinerU/Docling/OpenDataLoader/TCADP/PaddleOCR/Vision | |
_docx() | DOCX | DeepDOC + 大纲/页眉页脚/TOC 去除 |
_doc() | DOC | Apache Tika |
_spreadsheet() | XLSX/CSV | TCADP/DeepDOC ExcelParser |
_slides() | PPT/PPTX | TCADP/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 组件。
两种模式:
field_name == "toc" 时,调用 run_toc_from_text() 生成目录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_merge | TOC 去除,章节层级分块 |
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 等语音转文本 |
rag/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 构建流水线
generate_subgraph() 逐文档并行提取实体/关系merge_subgraph() 合并到全局 NetworkX 图resolve_entities() LLM 判断同名实体是否指代同一事物extract_community() Leiden 聚类 + LLM 生成社区摘要可恢复性:通过 Redis 阶段标记 (phase_markers.py) 支持中断后恢复。
三种提取策略:
general — 微软 GraphRAG 风格,LLM gleaning 迭代提取light — LightRAG 风格,简化 Promptner — spaCy NER + 共现关系,无需 LLMgraphrag/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:
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。
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.py | Elasticsearch | 全文/向量/融合搜索,Painless 脚本 Pagerank 调整,聚合 |
infinity_conn.py | Infinity | 多表搜索,字段名映射,DataFrame 字段提取 |
opensearch_conn.py | OpenSearch | KNN DSL 搜索,布尔查询 |
ob_conn.py | OceanBase | 全文+向量混合搜索(融合退化为独立),完整 SQL 模式定义 |
utils/redis_conn.py — Redis 工具 核心基础设施类:RedisDB (单例) — 完整 Redis 操作封装
功能:KV 操作、集合、有序集合、自增 ID、Lua 脚本事务、Stream 消息队列(生产者/消费者)、未确认消息重处理。
类:RedisDistributedLock — 分布式锁(delete-if-equal 安全释放)
rag/advanced_rag/tree_structured_query_decomposition_retrieval.py DeepResearcher类:TreeStructuredQueryDecompositionRetrieval (别名 DeepResearcher)
实现树形查询分解深度研究:
关键依赖:rag.prompts.generator (充分性检查/多查询生成), rag.utils.tavily_conn, rag.nlp.search
rag/svr/svr/task_executor.py — 任务执行器 核心Redis Stream 消费者,处理所有文档摄入任务:
# 并发控制信号量
task_limiter # 任务级并发限制
chunk_limiter # 分块并发限制
embed_limiter # 嵌入生成速率限制
minio_limiter # MinIO 并发限制
kg_limiter # 知识图谱并发限制
任务类型分发:
| 任务类型 | 处理函数 | 描述 |
|---|---|---|
| MEMORY | build_chunks() | 标准文档分块+嵌入 |
| DATAFLOW | run_dataflow() | 流水线执行 |
| RAPTOR | run_raptor_for_kb() | RAPTOR 摘要树 |
| GRAPHRAG | run_graphrag_for_kb() | 知识图谱构建 |
| MINDMAP | run_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。
同步机制:基于指纹的变更检测 + 定时轮询调度。
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 → 文档摄入管道 |