跳到主要内容

入库写入线 — 分块、分词、向量、索引

30 秒导读:01 章 把一份 PDF/DOCX 变成了一串 chunk(文本片段)。本章讲这些 chunk 如何变成可检索的索引记录:后台 worker 从 Redis 队列取任务 → 切块 → 中文分词 + 算词权重 → 编码成向量 → 写字段名有讲究的记录 → 批量灌进文档库(Elasticsearch / Infinity / OpenSearch / OceanBase)。范围止于「记录已落库、可被检索」;查询、打分、重排是第 03 章


1. 这是什么(零基础也能懂)

一句话定义: 这是 RAGFlow 的写入侧后台流水线——一个独立的 Python worker 进程,把「一段文本」加工成「文档库里的一行可搜索记录」。

为什么需要它。 检索要快,靠的是提前把每个 chunk 处理好、建好索引。查询时才现算就太慢了。所以入库阶段要一次性把三样东西都算出来、和文本一起存下:

要存的东西干什么用服务于
分词后的 token 串(content_ltks 等)全文/关键词匹配(BM25 打分的基座)第 03 章的文本召回
稠密向量(q_1024_vec 等)语义相似匹配(向量召回)第 03 章的向量召回
排序特征(pagerank_featag_feas)让重要/带标签的 chunk 排名靠前第 03 章的融合排序

它在整条链路的位置。 起点是「chunk 已由 chunker 产出」,终点是「记录已 insert 进 doc store」。中间不碰 deepdoc 的解析细节(那是第 01 章),也不碰查询(第 03 章)。

一句话直觉: 把入库想成给每个 chunk 办一张身份证——正面印原文,背面同时印上「分词指纹 + 向量指纹 + 优先级标记」,存进档案库后,检索时按任一种指纹都能快速找到它。

本节不出现代码细节。目标:让完全不懂的人知道「这条线在干嘛」。


2. 顶层全景(它大概怎么转)

2.1 一张图:从队列到落库

worker 是个永不退出的循环:抢一个并发额度 → 取一条任务 → 端到端处理 → ack。怎么读下图:从上到下是一条任务的一生,左侧是主循环,右侧是 do_handle_task 内部的四步加工。

Redis Stream (优先级队列) task_executor.py 进程
te.1.common (高) te.0.common (低)
│ ┌───────────────────────────┐
▼ │ main(): while 循环 │
① collect() 取一条任务 ◄─────────────────┤ 抢并发额度 task_limiter │
│ (消费者组 rag_flow_svr_task_broker)└───────────────────────────┘

② handle_task() → do_handle_task(task)

├─ Ⓐ build_chunks 切块 + 图片传 MinIO + 可选 LLM 增强
│ │ (chunker.chunk → cks)
│ ▼
├─ Ⓑ embedding 每个 chunk 编码成 q_<dim>_vec
│ │
│ ▼
├─ Ⓒ (分词已在 chunker 内做, 见 3.3) content_ltks / title_tks ...
│ │
│ ▼
└─ Ⓓ insert_chunks 分批 insert 进 doc store ──► ES / Infinity / ...


③ redis_msg.ack() + 更新进度/计数

2.2 部件一句话职责

部件干什么在哪(符号)
主循环无限抢额度、起协程处理任务rag/svr/task_executor.py:1851 main
取任务从两条优先级流取一条,查 DB 补全任务体rag/svr/task_executor.py:199 collect
任务分派task_type 走 dataflow/raptor/graphrag/普通 分支rag/svr/task_executor.py:1359 do_handle_task
切块chunker.chunk、传图、可选关键词/问题/元数据/标签rag/svr/task_executor.py:272 build_chunks
中文分词把文本切成 token 串rag/nlp/rag_tokenizer.py:18 RagTokenizer
词权重给每个 token 算 term/idf 权重rag/nlp/term_weight.py:27 Dealer
向量化批量编码成稠密向量rag/svr/task_executor.py:680 embedding
落库分批写多后端文档库rag/svr/task_executor.py:1243 insert_chunks
库抽象统一 ES/Infinity/OpenSearch/OceanBase 的接口common/doc_store/doc_store_base.py:143 DocStoreConnection
可编排入库新版 DAG(File→Parser→Chunker→Tokenizer→Extractor)rag/flow/pipeline.py:28 Pipeline

2.3 两条并存的入库路径(重要)

RAGFlow 现在同一个 worker 里并存两套入库实现,靠环境变量 TE_RUN_MODE 切换(rag/svr/task_executor.py:1699):

TE_RUN_MODE走哪套说明
"0"(默认)TaskManager.run_refactored_task重构版,拆进 rag/svr/task_executor_refactor/
"1"do_handle_task + dry_run_task 对比双跑对账,给迁移做校验
其它do_handle_task原始版,单文件里最完整可读的一条

本章讲原始版 do_handle_task 这条——它把「切块→分词→向量→落库」的每一步都摊在 task_executor.py 里,读起来最清楚;重构版是它的等价搬迁。另有一条全新的可编排 DAG(rag/flow/pipeline.py),用于用户自定义数据流,见 3.6


3. 核心原理(逐个机制)

3.1 后台 worker 主循环:从队列取一条任务

要解决的小问题: 上传文档后不能让 API 进程同步去解析(太慢会阻塞)。所以 API 只往 Redis Stream 塞一条任务,由多个独立 worker 进程并行消费。

主循环长这样。 main() 是个 while not stop_event:先 await task_limiter.acquire() 抢一个并发额度,再起一个协程去 handle_task()。额度用完就阻塞,天然做了背压(rag/svr/task_executor.py:1894-1898)。

取任务分两步走优先级。 collect() 先问「有没有我这个消费者没 ack 完的旧任务」(get_unacked_iterator),没有再按 高优先级流在前 依次消费(rag/svr/task_executor.py:207-215):

svr_queue_names = ["te.1.common", "te.0.common"] # 高→低, common/settings.py:159

▼ 逐个 queue_consumer, 命中即停
拿到 redis_msg → msg["id"] → TaskService.get_task(id) 从 DB 取完整任务体
  • 队列名由 get_svr_queue_names("common") 生成,固定返回 ["te.1.common","te.0.common"](common/settings.py:157)。
  • 消费者组名是常量 SVR_CONSUMER_GROUP_NAME = "rag_flow_svr_task_broker"(common/constants.py:261),多 worker 靠消费者组分摊负载。
  • 取到消息只是拿到 id;真正的任务字段(parser 配置、页码、语言、embedding 模型 id 等)要再查 DB(collectTaskService.get_task,rag/svr/task_executor.py:240)。

处理完必须 ack。 handle_task() 无论成功、取消、异常,最后都会 redis_msg.ack()(rag/svr/task_executor.py:1749),否则任务会被当成「未完成」重投。

坑: worker 还跑一个 report_status() 心跳协程,定期往 Redis 写自己的存活时间戳,并清理超时(默认 120s)的僵尸 worker(rag/svr/task_executor.py:1763WORKER_HEARTBEAT_TIMEOUT)。这让「某个 worker 崩了它的任务不会永远卡住」。


3.2 build_chunks:切块、传图、可选 LLM 增强

要解决的小问题: 把「文件二进制」变成「一串带字段的 chunk 字典」,并顺手把图片、关键词、问题、元数据、标签都备齐。

主线五步。 build_chunks(task, progress_callback)(rag/svr/task_executor.py:272):

① 取二进制 File2DocumentService.get_storage_address → STORAGE_IMPL.get (:283)
② 组 chunk_config 从 parser_config 读切块参数 (:313)
③ chunker.chunk 按 parser_id 分派到对应 chunker, 产出 cks (:329)
④ 传图 + 定 id 每个 chunk 深拷 doc 模板, 算 id, 图片上传 MinIO (:372 upload_to_minio)
⑤ 可选 LLM 增强 auto_keywords / auto_questions / metadata / tag (:413..614)

① chunker 由 parser_id 选。 一张工厂表把 parser_id 映射到具体 chunker 模块(rag/svr/task_executor.py:114 FACTORY),比如 naive→通用切块、paper/book/table/qa 各有专门策略。真正的 chunker.chunk(...) 调用在 rag/svr/task_executor.py:329,放进线程池跑(它是同步 CPU 活)。

② 切块参数三件套(chunk_config)。 这是「一份文档切多细」的旋钮,从 parser_config 取,带默认值(rag/svr/task_executor.py:313-322):

参数默认含义
chunk_token_num128每个 chunk 的目标 token 数
overlapped_percent0相邻 chunk 的重叠比例(经 normalize_overlapped_percent 归一)
delimiter"\n!?。;!?"优先在这些分隔符处切,避免切碎句子

③ chunk 出厂时已带全文字段。 通用 chunker 在切块时就调了分词器:每个 chunk 落成 content_with_weight(原文)+ content_ltks/content_sm_ltks(分词串),文档级还带 docnm_kwd(文件名)与 title_tks(标题分词)。见 rag/nlp/__init__.py:268 tokenizerag/app/naive.py:865所以分词发生在切块内,不是入库时另起一步——这点容易看错。

④ chunk 的 id 是内容哈希。 upload_to_minio 里用 xxhash.xxh64(content_with_weight + doc_id) 当 id(rag/svr/task_executor.py:376)。同文档同内容 → 同 id → 天然幂等去重,重复入库会覆盖而非堆叠。带图的 chunk 走 image2id 把图片存进 MinIO 并写 img_id(:389)。

⑤ 四个可选的 LLM 增强(默认关)。 每个都缓存(get_llm_cache/set_llm_cache)避免重复烧钱,产物挂到 chunk 上供检索用:

开关(parser_config)产出字段干什么位置
auto_keywordsimportant_kwd/important_tks抽关键词,增强关键词召回:413
auto_questionsquestion_kwd/question_tks生成「这段能回答的问题」,做问题匹配:450
enable_metadata文档级 metadataLLM 抽结构化元数据,供过滤:486
tag_kb_idstag_feas(TAG_FLD)打标签,写成 rank_features 影响排序:548

教学示例(示意,非源码):关键词增强的核心动作

# 从一个 chunk 抽 topn 个关键词, 再把关键词也分词, 双双挂回 chunk
cached = keyword_extraction(chat_mdl, chunk["content_with_weight"], topn) # LLM 抽词
chunk["important_kwd"] = [k for k in re.split(r"[,,;;、]+", cached) if k.strip()]
chunk["important_tks"] = rag_tokenizer.tokenize(" ".join(chunk["important_kwd"]))
# 重点看: 抽出来的词也要过分词器, 才能进全文索引参与匹配

对应真实实现 rag/svr/task_executor.py:426-430 doc_keyword_extraction


3.3 中文分词与词权重:BM25 打分的基座

要解决的小问题: 中文没有空格分词,"知识库检索" 得先切成 知识 库 检索 这样的 token,全文索引和 BM25(一种经典关键词打分算法,按词频/逆文档频率给匹配打分)才能工作。这一步的产物,是第 03 章文本召回的全部基础。

分词器是个薄封装。 rag/nlp/rag_tokenizer.py:18RagTokenizer 继承自 infinity.rag_tokenizer.RagTokenizer(底层 C++ 实现),对外暴露 tokenize(粗粒度)、fine_grained_tokenize(细粒度)、tag(词性)、freq(词频)。

一个关键分支:Infinity 后端不预分词。DOC_ENGINE_INFINITY 为真时,tokenize 直接原样返回文本(rag/nlp/rag_tokenizer.py:22),因为 Infinity 引擎自己在库内分词。所以「是否预分词」取决于文档库后端。

词权重给 token 定「谁更重要」。 rag/nlp/term_weight.py:27Dealer 把一串 token 变成 [(token, 归一化权重), ...]。核心是 weights()(:164),权重 = 两种 idf 的加权 × 命名实体系数 × 词性系数:

权重(t) = (0.3·idf(freq) + 0.7·idf(df)) × ner(t) × postag(t)
│词频逆频 │文档逆频 │实体类型 │词性
└── term_weight.py:225 idf ──┘ :170 ner :181 postag
最后整串归一化: 每个权重 / 权重之和 (term_weight.py:246-247)

各因子的直觉(都在 weights 的内嵌函数里):

因子直觉例子(源码常量)
ner(t)是机构/地名/股票这类实体就加权corp/loca/sch/stock → 3(:177)
postag(t)名词、专名重,代词/连词/副词轻ns/nt → 3,n → 2,r/c/d → 0.3(:183-191)
freq(t)/df(t)越稀有的词区分度越高、权重越大长词拆细粒度取最小值(:202:218)

为什么重要: 这套权重在入库时不直接落库,而是查询时对 query 做同样处理来做查询扩展与加权(第 03 章的 Dealer 会复用它)。入库侧真正落库的是分词串本身(*_ltks/*_tks),权重逻辑是「查询和文档共用的同一把尺子」。


3.4 向量化:把 content 编码成 q__vec

要解决的小问题: 语义检索需要把每个 chunk 变成一个稠密向量,查询时用向量相似度找「意思相近」的片段。

主线。 embedding(docs, mdl, parser_config, callback)(rag/svr/task_executor.py:680):

  1. 取要编码的文本。 优先用 question_kwd(若做了 auto_questions),否则用 content_with_weight;并把表格标签 <table><td>... 去掉再编码(:686-693)。
  2. 标题单独编码一次、再和正文加权融合。 标题只编码一份、平铺到所有 chunk(:697-699),然后按 filename_embd_weight(默认 0.1)融合:vects = 0.1·标题向量 + 0.9·正文向量(:717-719)。好处:文件名信息渗进每个 chunk 的向量,但只占一成。
  3. 分批编码。EMBEDDING_BATCH_SIZE 一批,超长文本先 truncate 到模型上限再编码(:702-712)。
  4. 按维度命名字段写回。 关键一行:
# rag/svr/task_executor.py:728
d["q_%d_vec" % len(v)] = v # 1024 维 → 字段名 "q_1024_vec"

为什么字段名带维度? 因为文档库的向量字段是按维度动态映射的(见 3.5 的 *_1024_vec 模板)。字段名把维度编进去,一个索引就能容纳不同 embedding 模型的向量;也让检索时能精确挑对字段。

坑(inferred 度低,代码可见): RAPTOR / dataflow 分支各自还有一份独立的向量化循环(rag/svr/task_executor.py:806 / :1102),不复用 embedding()——因为它们的输入结构不同(summary、pipeline chunk)。


3.5 落库:多后端 doc store 抽象 + 字段命名约定

要解决的小问题: RAGFlow 要支持四种文档库(ES / Infinity / OpenSearch / OceanBase),上层入库代码不能为每种写一遍。

统一接口。 一个抽象基类 DocStoreConnection(common/doc_store/doc_store_base.py:143)定义了 create_idx / insert / search / delete / update 等抽象方法;每种后端一个实现(rag/utils/es_conn.pyinfinity_conn.pyopensearch_conn.pyob_conn.py),运行时 settings.docStoreConn 指向选中的那个。上层只认接口。

建索引在入库前。 do_handle_task 先探测向量维度(用 "ok" 试编码一次),再 init_kb(task, vector_size) 建索引(rag/svr/task_executor.py:1405-1413);init_kb 转调 docStoreConn.create_idx(:673)。ES 实现直接套用 conf/mapping.json 的 settings/mappings(common/doc_store/es_conn_base.py:128)。

批量 insert。 insert_chunksDOC_BULK_SIZE 分批调 docStoreConn.insert(rag/svr/task_executor.py:1288-1294);ES 实现用 bulk API,以 chunk 的 id_id 保证唯一(rag/utils/es_conn.py:307)。每批后更新 TaskService.update_chunk_ids 记录进度,任务被取消就回滚(:1300-1321)。

精髓:字段类型由「字段名后缀」决定。 ES 索引不写死每个字段,而是用 dynamic_templates 按后缀匹配(conf/mapping.json)。这解释了为什么全代码库的字段名都长成 xxx_kwdxxx_ltks 这样:

字段名后缀映射成的 ES 类型干什么例子字段
*_ltkstext, whitespace 分析器粗粒度分词全文content_ltks
*_tkstext, scripted_sim 相似度分词全文(参与 BM25 打分)title_tksimportant_tks
*_kwd / *_idkeyword精确匹配/过滤docnm_kwdimportant_kwd
*_with_weighttext, index:false只存不建索引(存原文)content_with_weight
*_int / *_fltinteger / float数值(位置、时间戳)page_num_intcreate_timestamp_flt
*_fearank_feature单值排序特征pagerank_fea
*_feasrank_features多值排序特征tag_feas
*_<dim>_vecdense_vector, cosine稠密向量q_1024_vec

来源:conf/mapping.jsondynamic_templates(*_tks/*_ltks/*_kwd/*_with_weight/*_fea/*_feas/*_512_vec~*_1536_vec 逐条)。

两个 rank feature 字段的写入时机。 它们是第 03 章排序的输入,在入库侧就写好:

  • PAGERANK_FLD = "pagerank_fea"(common/constants.py:259):若知识库配了 pagerank,build_chunks 把它塞进每个 chunk 的 doc 模板(rag/svr/task_executor.py:367-368)。
  • TAG_FLD = "tag_feas"(common/constants.py:262):auto-tag 分支把标签字典写进 d[TAG_FLD](:571:597),落库成 rank_features,检索时命中标签的 chunk 得到 boost。

母块(mother chunk)机制: insert_chunks 开头会把带 mom/mom_with_weight 的 chunk 抽出「父片段」单独入库一份(rag/svr/task_executor.py:1254-1273),父块只保留少数字段、available_int=0(不直接召回)。这是「小块命中、返回大块上下文」的存储基础。


3.6 新版可编排入库 DAG(pipeline)

要解决的小问题: 上面 do_handle_task写死的固定流水线。用户想自定义「怎么解析、怎么切、要不要抽取」,就需要一个可编排的图

Pipeline 是一张有向图。 rag/flow/pipeline.py:28Pipeline 继承自 agent.canvas.Graph,用 DSL(JSON)描述节点和连线。典型链路:

File ──► Parser ──► Chunker ──► Tokenizer ──► Extractor
取二进制 分派解析 切块 分词+向量 LLM 抽取

执行靠拓扑推进。 run()(rag/flow/pipeline.py:117)从 File 起步,每步把上游节点的 output() 作为下游 invoke(**kwargs) 的入参,沿 path 一路推进,出错即停并回写进度(:143-166)。

节点间用 Pydantic schema 校验数据契约。 每个节点声明一个 *FromUpstream 模型,验上游给的数据合法才干活。以 Tokenizer 为例(rag/flow/tokenizer/schema.py:20 TokenizerFromUpstream):它声明 output_format 可为 chunks/markdown/text/html/json,并按格式校验对应 payload 是否存在(:38 _check_payloads)。这让「上游产出什么、下游期待什么」成为显式契约,而非隐式字典约定。

Tokenizer 节点的产物和传统路径一致。 它同样算 content_ltks/content_sm_ltks/title_tks 并写 q_%d_vec(rag/flow/tokenizer/tokenizer.py:114:146-150),search_method 决定做全文、向量、还是两者都做(:129)。

两条路径的关系。

维度传统 do_handle_task可编排 Pipeline
流程固定四步,写死用户 DSL 描述的 DAG
触发task_type = 普通解析task_typedataflow 开头
入口rag/svr/task_executor.py:1539run_dataflowPipeline.run(:733:752)
落库insert_chunks(同一函数)insert_chunks(同一函数,:882)
数据契约隐式字典字段显式 *FromUpstream Pydantic schema

关键:两条路殊途同归。 dataflow 跑完 pipeline 拿到 chunks,补齐 id/时间戳/分词字段后,照样调同一个 insert_chunks 落库(rag/svr/task_executor.py:882)。所以无论走哪条路,最终写进 doc store 的记录形态是一致的,第 03 章的检索无需区分来源。


4. 一条 chunk 落库时的字段全景

把前面各步的产物汇总——一条普通文本 chunk 入库时,记录大致长这样(字段→来源步骤):

字段内容由哪步写
idxxhash(content+doc_id)build_chunks(:376)
doc_id / kb_id归属文档/知识库build_chunks(:366)
docnm_kwd / title_tks文件名 / 标题分词chunker(naive.py:865)
content_with_weight原文(只存不索引)chunker(nlp/__init__.py:270)
content_ltks / content_sm_ltks正文粗/细粒度分词chunker(nlp/__init__.py:272)
important_kwd / important_tks关键词(可选)build_chunks(:429)
question_kwd / question_tks问题(可选)build_chunks(:466)
q_<dim>_vec稠密向量embedding(:728)
pagerank_feapagerank 排序特征(可选)build_chunks(:368)
tag_feas标签 rank_features(可选)build_chunks(:597)
create_timestamp_flt / img_id时间戳 / 图片引用build_chunks(:378:389)

5. 巧妙之处(可借鉴的技术)

  1. 字段名即 schema。 用 ES dynamic_templates 把「后缀 → 类型」约定成规则(conf/mapping.json),字段名 content_ltks/pagerank_fea/q_1024_vec 自带类型信息。加新字段不用改 mapping,只要起对名字。代价是命名纪律必须全库统一。
  2. 内容哈希当 id,天然幂等。 xxhash(content+doc_id) 做 id(:376),重复解析同一文档不会产生重复记录,重跑安全。
  3. 向量字段名编进维度。 q_%d_vec(:728)让一个索引兼容多种 embedding 模型的向量维度,换模型不冲突。
  4. 标题向量低权融合。 0.1·标题 + 0.9·正文(:717)把文件名信号渗进每个 chunk 又不喧宾夺主。
  5. LLM 增强全缓存。 keywords/questions/metadata/tags 都过 get_llm_cache(:420 等),重跑不重复烧 token。
  6. 两条入库路殊途同归到 insert_chunks 固定流水线和可编排 DAG 最终落库形态一致,下游检索无需分支。

6. 边界与局限(诚实)

  • 本章止于「落库可检索」。 向量相似度、BM25、rerank、融合排序都在第 03 章;RAPTOR/GraphRAG 的进阶结构在第 04 章
  • 默认走重构版而非 do_handle_task TE_RUN_MODE="0" 时实际执行 TaskManager.run_refactored_task(:1713),do_handle_task 是逻辑等价的可读参照。本章按后者讲解。
  • Infinity 后端不预分词。 tokenizeDOC_ENGINE_INFINITY 下原样返回(rag/nlp/rag_tokenizer.py:22),分词逻辑下沉到库内,所以 *_ltks 字段在不同后端语义不同。
  • term_weight 的权重不直接落库。 它是查询/文档共用的打分尺子,入库侧落的是分词串;别误以为每个 token 的权重被存进了索引。
  • 母块机制依赖 chunker 产出 mom 字段。 没有该字段就不会生成父片段(:1256-1259),父子召回也就退化。

7. 横向对比(同 shelf 兄弟)

项目入库侧取舍对照
RAGFlow(本章)重字段工程:一条 chunk 同时备好分词/向量/rank feature,查询时几乎不现算
通用向量库 RAG多数只存「向量 + 原文」,BM25 与排序特征靠外部引擎RAGFlow 把 BM25 基座(分词串)也预存进同一记录

RAGFlow 的特点是把「文本索引 + 向量索引 + 排序特征」压进同一条记录、同一个库,以便第 03 章做「向量 + BM25 + rank feature」的一次性融合召回。跨库检索原理见 shelf 总库 ai-agent-reference 的 rag-retrieval 主题。


8. 代码地图(导航索引)

主题文件路径符号
主循环rag/svr/task_executor.pymain handle_task task_manager
取任务rag/svr/task_executor.pycollect
任务分派rag/svr/task_executor.pydo_handle_task
切块+增强rag/svr/task_executor.pybuild_chunks upload_to_minio doc_keyword_extraction
建索引rag/svr/task_executor.pyinit_kb
向量化rag/svr/task_executor.pyembedding
落库rag/svr/task_executor.pyinsert_chunks
队列名common/settings.pyget_svr_queue_name get_svr_queue_names
分词器rag/nlp/rag_tokenizer.pyRagTokenizer tokenize fine_grained_tokenize
词权重rag/nlp/term_weight.pyDealer weights ner postag freq df idf
切块内分词rag/nlp/__init__.pytokenize tokenize_chunks add_positions
库抽象common/doc_store/doc_store_base.pyDocStoreConnection create_idx insert
ES 实现rag/utils/es_conn.py · common/doc_store/es_conn_base.pyinsert create_idx
字段映射conf/mapping.jsondynamic_templates
rank 特征常量common/constants.pyPAGERANK_FLD TAG_FLD SVR_CONSUMER_GROUP_NAME
可编排 DAGrag/flow/pipeline.pyPipeline run
DAG 数据契约rag/flow/tokenizer/schema.py · tokenizer.pyTokenizerFromUpstream
dataflow 入口rag/svr/task_executor.pyrun_dataflow