RAG 总览与索引:一份 PDF 怎么变成"可检索的知识"
从今天起进入 RAG(检索增强生成)。核心矛盾是:大模型没读过你公司的内部文档,你又不可能把 500 页手册整个塞进 Prompt。RAG 的答案是——先把文档切成小片段、建好索引存起来(离线索引),提问时只捞出最相关的几片喂给模型(在线检索)。今天专攻前半段"建索引",主角是 core/indexing_runner.py。弄清三件事:①一份文档从原始文件到可检索片段,要过哪五步;②文本怎么被"切分"成大小合适的块;③每一块怎么被向量化、写进向量库。
痛点:500 页手册,怎么让模型"读得懂又用得上"?
core/indexing_runner.py 就是这条离线流水线的总指挥。今天看它,Day14 看在线检索那半段。本质:一条流水线,五个工位
| 工位 | 方法 | 干什么 |
|---|---|---|
| ① 抽取 | _extract | 把 PDF/Word/网页等原始文件解析成纯文本 Document |
| ② 转换 | _transform indexing_runner.py:780 | 清洗 + 切分成小块 + 给每块打指纹(这一步含切分) |
| ③ 存片段 | _load_segments indexing_runner.py:816 | 把切好的片段写进业务库(document_segments 表) |
| ④ 建索引 | _load indexing_runner.py:573 | 向量化 + 写进向量库/关键词索引 |
| —— 贯穿 —— | IndexProcessor | 按 doc_form(普通/父子/QA)选不同处理策略 |
IndexingRunner 是"车间主任",它不亲自干活,而是按顺序喊每个工位干活。有意思的是它先根据文档类型(普通段落 / 父子分层 / 问答对)挑一个 IndexProcessor(工艺流程单),后面抽取、转换、建索引全委托给这张流程单——同一条主流程,不同工艺可插拔,这是策略模式的经典用法。run():把五步串起来
总入口 IndexingRunner.run(core/indexing_runner.py:68)——一眼看清整条流水线:
# core/indexing_runner.py:68
def run(self, dataset_documents: list[DatasetDocument]):
for dataset_document in dataset_documents:
try:
dataset = db.session.get(Dataset, requeried_document.dataset_id) # 找到所属知识库
processing_rule = db.session.scalar(...) # 取处理规则(切分参数等)
index_type = requeried_document.doc_form
index_processor = IndexProcessorFactory(index_type).init_index_processor() # ① 按类型选工艺
text_docs = self._extract(index_processor, requeried_document, processing_rule.to_dict()) # ② 抽取
documents = self._transform(index_processor, dataset, text_docs, ...) # ③ 清洗+切分
self._load_segments(dataset, requeried_document, documents) # ④ 存片段到业务库
self._load(index_processor=index_processor, dataset=dataset, # ⑤ 向量化+建索引
dataset_document=requeried_document, documents=documents)
except DocumentIsPausedError:
raise DocumentIsPausedError(...) # 用户暂停:原样抛出
except ProviderTokenNotInitError as e:
self._handle_indexing_error(document_id, e) # 没配 embedding key:记错
except Exception as e:
self._handle_indexing_error(document_id, e) # 兜底:把错写进文档状态
IndexProcessorFactory(index_type)★先按 doc_form(普通段落 / 父子分层 / QA 问答)造一个对应的处理器。后面 extract/transform/load 都调它——主流程只有一份,工艺可换。五个方法顺序调用抽取 → 转换 → 存片段 → 建索引,一条直线。读源码时抓住这四行,整个索引流程就懂了七成。三层 except★暂停原样抛(让上层知道是主动停的,不是真错);没配 embedding key 单独接(这是最常见的用户配置错误);其余 兜底 把异常写进文档的 error 字段。_handle_indexing_error 会把文档状态标成 error——所以你在界面看到"索引失败"红字,就是这里写的。run_in_splitting_status(indexing_runner.py:128)和 run_in_indexing_status(indexing_runner.py:199)是"断点续跑":如果文档卡在"切分中"或"索引中"状态(比如上次跑到一半进程挂了),从对应工位重新开始而不是从头。它俩会先删掉上次残留的 segment 再重来——保证幂等,不会重复入库。切分:_get_splitter 怎么把长文断成块
切分是索引质量的命门。切分器由 BaseIndexProcessor._get_splitter(core/rag/index_processor/index_processor_base.py:100)按规则挑:
# core/rag/index_processor/index_processor_base.py:100
def _get_splitter(self, processing_rule_mode, max_tokens, chunk_overlap, separator,
embedding_model_instance) -> TextSplitter:
if processing_rule_mode in ["custom", "hierarchical"]: # 用户自定义规则
max_segmentation_tokens_length = dify_config.INDEXING_MAX_SEGMENTATION_TOKENS_LENGTH
if max_tokens < 50 or max_tokens > max_segmentation_tokens_length:
raise ValueError(f"Custom segment length should be between 50 and {max_segmentation_tokens_length}.")
character_splitter = FixedRecursiveCharacterTextSplitter.from_encoder(
chunk_size=max_tokens, chunk_overlap=chunk_overlap,
fixed_separator=separator,
separators=["\n\n", "。", ". ", " ", ""], # ★ 优先在段落/句号处断
embedding_model_instance=embedding_model_instance)
else: # 自动模式:用内置默认参数
character_splitter = EnhanceRecursiveCharacterTextSplitter.from_encoder(
chunk_size=DatasetProcessRule.AUTOMATIC_RULES["segmentation"]["max_tokens"],
chunk_overlap=DatasetProcessRule.AUTOMATIC_RULES["segmentation"]["chunk_overlap"],
separators=["\n\n", "。", ". ", " ", ""],
embedding_model_instance=embedding_model_instance)
return character_splitter
chunk_size / max_tokens每块最大多少 token。太大:一块塞太多内容,检索精度下降、还费上下文;太小:语义被切碎、丢上下文。所以校验 50 ≤ max_tokens ≤ 上限。separators 的顺序★["\n\n", "。", ". ", " ", ""] 是递归切分的关键:优先在"空行/段落"断,其次句号,再次空格,实在不行才逐字符切。这样尽量不把一句话拦腰截断,保住语义完整。chunk_overlap相邻块重叠一小段。为什么要重叠?防止"关键信息正好落在两块交界处"被割裂——重叠让边界信息在两块里都出现一次。from_encoder(embedding_model_instance=...)用 embedding 模型的分词器来数 token——因为最终是喂给这个模型,用它的口径数长度才准。transform:清洗 + 切分 + 打指纹
以最常见的"普通段落"工艺为例,看 ParagraphIndexProcessor.transform(core/rag/index_processor/processor/paragraph_index_processor.py:75):
# core/rag/index_processor/processor/paragraph_index_processor.py:75
def transform(self, documents, current_user=None, **kwargs) -> list[Document]:
process_rule = kwargs.get("process_rule")
rules = Rule.model_validate(...) # 解析切分规则
splitter = self._get_splitter(...) # ① 拿到 L04 的切分器
all_documents = []
for document in documents:
document_text = CleanProcessor.clean(document.page_content, kwargs.get("process_rule", {})) # ② 清洗
document.page_content = document_text
document_nodes = splitter.split_documents([document]) # ③ ★真正切分
split_documents = []
for document_node in document_nodes:
if document_node.page_content.strip():
doc_id = str(uuid.uuid4())
hash = helper.generate_text_hash(document_node.page_content) # ④ 内容指纹
document_node.metadata["doc_id"] = doc_id
document_node.metadata["doc_hash"] = hash
page_content = remove_leading_symbols(document_node.page_content).strip()
if len(page_content) > 0:
document_node.page_content = page_content
split_documents.append(document_node)
all_documents.extend(split_documents)
return all_documents
CleanProcessor.clean去杂质:多余空白、控制字符、按规则去 URL/邮箱等。垃圾进 = 垃圾出,清洗直接影响后面向量的质量。splitter.split_documents★调 L04 的切分器,把一篇长文切成一个个 document_node(片段)。这是 transform 的核心动作。doc_id + doc_hash★每片打两个标记:doc_id(唯一 id,向量库里的主键)和 doc_hash(内容指纹)。指纹的用处:内容没变就不用重新向量化,增量更新时靠它去重、省 embedding 调用。strip() 过滤空块切出来的空白片段直接丢弃——空块进向量库既浪费又污染检索结果。load:向量化 + 落库
片段切好后,ParagraphIndexProcessor.load(core/rag/index_processor/processor/paragraph_index_processor.py:125)负责把它们送进索引:
# core/rag/index_processor/processor/paragraph_index_processor.py:125
def load(self, dataset, documents, multimodal_documents=None, with_keywords=True, **kwargs):
if dataset.indexing_technique == IndexTechniqueType.HIGH_QUALITY: # 高质量:走向量
vector = Vector(dataset)
vector.create(documents) # ★ 向量化 + 写进向量数据库
if multimodal_documents and dataset.is_multimodal:
vector.create_multimodal(multimodal_documents)
with_keywords = False
if with_keywords: # 经济模式:走关键词索引(不花 embedding 钱)
keywords_list = kwargs.get("keywords_list")
keyword = Keyword(dataset)
if keywords_list and len(keywords_list) > 0:
keyword.add_texts(documents, keywords_list=keywords_list)
else:
keyword.add_texts(documents)
而真正调 embedding 模型、并发向量化的编排在 IndexingRunner._load(core/indexing_runner.py:573)里:
# core/indexing_runner.py:573
def _load(self, index_processor, dataset, dataset_document, documents):
embedding_model_instance = None
if dataset.indexing_technique == IndexTechniqueType.HIGH_QUALITY:
embedding_model_instance = self._get_model_manager(dataset.tenant_id).get_model_instance(
tenant_id=dataset.tenant_id, provider=dataset.embedding_model_provider,
model_type=ModelType.TEXT_EMBEDDING, model=dataset.embedding_model) # ★ 拿 embedding 模型
...
max_workers = 10
if dataset.indexing_technique == IndexTechniqueType.HIGH_QUALITY:
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
document_groups = [[] for _ in range(max_workers)]
for document in documents:
hash = helper.generate_text_hash(document.page_content)
group_index = int(hash, 16) % max_workers # ★ 按内容 hash 分组,避免同片撞车
document_groups[group_index].append(document)
HIGH_QUALITY vs ECONOMY★两种索引技术:高质量用 embedding 模型做向量(语义检索,效果好但要花钱调模型);经济只建关键词索引(免费,但只能字面匹配)。load 里用 with_keywords 开关区分。这直接对应你建知识库时选的那个选项。get_model_instance(TEXT_EMBEDDING)★还记得 Day05 的 ModelManager 吗?这里用它拿 embedding 模型实例——RAG 和模型运行时在这里接头。ThreadPoolExecutor + hash 分组★向量化是网络 IO 密集(要一批批调 embedding API),所以用 10 个线程并发。按 hash % 10 把片段分到不同组,源码注释点明:避免多线程处理"同一份文档"引发数据库插入死锁。vector.create(documents)把每片的向量写进向量库(Dify 支持多种:Qdrant/Weaviate/PGVector…,由 Vector 抽象统一)。Day14 检索就是从这里查。状态机 + 今日小结
存片段那一步 _load_segments(core/indexing_runner.py:816)除了写库,还推进文档的状态机:
# core/indexing_runner.py:816
def _load_segments(self, dataset, dataset_document, documents):
doc_store = DatasetDocumentStore(dataset=dataset, user_id=dataset_document.created_by,
document_id=dataset_document.id)
doc_store.add_documents(docs=documents, save_child=...) # 片段写进 document_segments 表
self._update_document_index_status( # ★ 文档状态 → INDEXING
document_id=dataset_document.id,
after_indexing_status=IndexingStatus.INDEXING,
extra_update_params={..., DatasetDocument.word_count: sum(len(d.page_content) for d in documents)})
self._update_segments_by_document( # ★ 每个片段状态 → INDEXING
dataset_document_id=dataset_document.id,
update_params={DocumentSegment.status: SegmentStatus.INDEXING, ...})
run() 按 doc_form=paragraph 选处理器 → _extract 解析出约 5 万字纯文本 → _transform 清洗后按 max_tokens=500、overlap=50 切成约 120 个片段,每片打上 doc_id 和 doc_hash → _load_segments 把 120 条写进 document_segments 表,文档状态转 indexing → _load 用 10 线程并发调 embedding 模型,把每片转成 1536 维向量写进向量库 → 全部完成,状态转 completed。此后用户问"年假几天",就能在毫秒级从这 120 片里捞出讲年假的那两片(Day14 的活)。👶 小白:向量到底是什么?为什么"向量接近"就等于"意思相近"?
👨🏫 老师:向量就是一串数字(比如 1536 个小数),是 embedding 模型给一段文字算出的"语义坐标"。模型在海量语料上学会了:意思相近的文字,坐标也靠得近。比如"报销流程"和"费用如何申请"字面完全不同,但它们的向量在空间里挨得很近;而"报销流程"和"食堂菜单"的向量离得很远。所以"找语义最相关的片段"就变成了数学问题——在向量空间里找离问题向量最近的几个点。这也是为什么它比传统关键词搜索强:能理解"换个说法的同一个意思"。
🧠 今天你应该能回答
- RAG 为什么要建索引?(文档塞不进 Prompt,靠离线切分+向量化+检索)
- 索引流水线哪五步?(抽取→清洗→切分→存片段→向量化建索引)
- 切分为什么用递归 separators?(尽量不拦腰截断句子,保住语义)
- chunk_overlap 干嘛的?(相邻块重叠,防边界信息被割裂)
- doc_hash 有什么用?(内容指纹,增量更新去重、不重复向量化)
- 高质量 vs 经济模式区别?(向量语义检索 vs 关键词字面索引)
✋ 10 分钟动手
cd /Users/bitmart/work/codes/github/AI_WORK/dify/api
# 1. 主流水线
sed -n '68,127p' core/indexing_runner.py # run() 五步
sed -n '780,814p' core/indexing_runner.py # _transform(挑 embedding 模型)
sed -n '573,620p' core/indexing_runner.py # _load(并发向量化)
# 2. 切分与转换
sed -n '100,137p' core/rag/index_processor/index_processor_base.py # _get_splitter
sed -n '75,145p' core/rag/index_processor/processor/paragraph_index_processor.py # transform + load
# 3. 有哪些处理工艺
ls core/rag/index_processor/processor/
core/rag/datasource/retrieval_service.py。