Day 13 / 共 20 天 · 阶段4 RAG 知识库

RAG 总览与索引:一份 PDF 怎么变成"可检索的知识"

从今天起进入 RAG(检索增强生成)。核心矛盾是:大模型没读过你公司的内部文档,你又不可能把 500 页手册整个塞进 Prompt。RAG 的答案是——先把文档切成小片段、建好索引存起来(离线索引),提问时只捞出最相关的几片喂给模型(在线检索)。今天专攻前半段"建索引",主角是 core/indexing_runner.py。弄清三件事:①一份文档从原始文件到可检索片段,要过哪五步;②文本怎么被"切分"成大小合适的块;③每一块怎么被向量化、写进向量库。

📍 你在 20 天里的位置(阶段4:RAG 知识库)
S2 模型运行时 S3 工作流引擎 D13 RAG 与索引 D14 检索 retrieval D15 知识库数据流 S5 工具/Agent S6 收官
💡 先用两个类比兜住今天 类比一:建索引像给一本厚书做"读书卡片 + 目录"。你不会让人临时通读全书找答案,而是提前把书拆成一张张卡片(切分),每张卡片记下主题标签方便快速定位(向量化)。以后要查什么,翻卡片就行。类比二:整条索引流水线像食品加工厂——原料进厂(抽取原文)→ 清洗去杂质(清洗)→ 切成小块(切分)→ 贴上营养标签装袋(向量化 + 元数据)→ 入库上架(写向量库)。每一步都是独立工位,出问题能精确定位到哪一站。
L01

痛点:500 页手册,怎么让模型"读得懂又用得上"?

🤔 痛点你想做一个"公司制度问答助手",手里有几十份 PDF/Word/Markdown。直接把全文塞进 Prompt?不行——上下文窗口装不下、还贵得离谱。让模型"记住"?模型是预训练好的,不认识你的私有文档。而且用户问"报销流程是什么",你需要精准找到手册里讲报销那两三段,而不是把整本手册都发过去。问题是:怎么把海量文档变成"能按语义快速捞出相关片段"的东西?
💡 本质:把"全文检索"升级成"语义检索",且离线预处理RAG 的核心是把重活挪到离线:提前把文档切成小片段(chunk),把每片用 embedding 模型转成一个向量(一串能代表语义的数字),存进向量数据库。提问时,把问题也转成向量,去库里找"向量最接近的几片"——这就是"语义相似"。core/indexing_runner.py 就是这条离线流水线的总指挥。今天看它,Day14 看在线检索那半段。
L02

本质:一条流水线,五个工位

工位方法干什么
① 抽取_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(工艺流程单),后面抽取、转换、建索引全委托给这张流程单——同一条主流程,不同工艺可插拔,这是策略模式的经典用法。
L03

run():把五步串起来

总入口 IndexingRunner.runcore/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_statusindexing_runner.py:128)和 run_in_indexing_statusindexing_runner.py:199)是"断点续跑":如果文档卡在"切分中"或"索引中"状态(比如上次跑到一半进程挂了),从对应工位重新开始而不是从头。它俩会先删掉上次残留的 segment 再重来——保证幂等,不会重复入库。
L04

切分:_get_splitter 怎么把长文断成块

切分是索引质量的命门。切分器由 BaseIndexProcessor._get_splittercore/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——因为最终是喂给这个模型,用它的口径数长度才准。
💡 设计取舍:块大小的两难切分参数是 RAG 里最需要调的旋钮。块大 → 召回时上下文全,但一块里混入无关内容会稀释相关性、也更费 token;块小 → 定位精准,但容易丢上下文(比如"它"指代的主语在上一块)。Dify 给了自动模式(内置默认)和自定义模式(你自己填 max_tokens / overlap / separator),外加 overlap 来缓解边界问题。没有万能参数,只有"针对你的文档试出来"的参数
L05

transform:清洗 + 切分 + 打指纹

以最常见的"普通段落"工艺为例,看 ParagraphIndexProcessor.transformcore/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() 过滤空块切出来的空白片段直接丢弃——空块进向量库既浪费又污染检索结果。
离线索引流水线:原始文件 → 可检索片段 ① 抽取 PDF→文本 ② 清洗 去杂质 ③ 切分 递归断句 ④ 打指纹 doc_id/hash ⑤ 向量化 写向量库 ① _extract | ②③④ _transform | 存片段 _load_segments | ⑤ _load chunk = { page_content: "报销需在30天内…", metadata:{doc_id, doc_hash}, vector:[0.02, -0.11, …] }
图注:五个工位各司其职,产物是一条条"带向量和指纹的片段",等着被检索。
L06

load:向量化 + 落库

片段切好后,ParagraphIndexProcessor.loadcore/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._loadcore/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 检索就是从这里查。
L07

状态机 + 今日小结

存片段那一步 _load_segmentscore/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, ...})
文档在库里有一条状态链:waiting → parsing → cleaning → splitting → indexing → completed(失败则 error)。每过一个工位就更新一次状态。好处:界面能实时显示进度条、断了能从状态断点续跑、出错能定位卡在哪一站。这就是为什么上传大文档时你能看到"解析中/索引中"的原因。
📝 真实值:上传《员工手册.pdf》建索引 你把 30 页的《员工手册.pdf》拖进"高质量"知识库 → run()doc_form=paragraph 选处理器 → _extract 解析出约 5 万字纯文本 → _transform 清洗后按 max_tokens=500、overlap=50 切成约 120 个片段,每片打上 doc_iddoc_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/
明日预告 · Day 14:索引建好了,明天看在线的另一半——检索 retrieval。用户一句提问,Dify 怎么并发跑向量检索/全文检索/关键词检索,再用 rerank(重排)把最相关的片段顶到前面?主角是 core/rag/datasource/retrieval_service.py
← Day 12 事件与流式 Day 14 · 检索 retrieval →