🔍 RAGFlow QA 解析→检索 全流程源码分析

从文件上传到 LLM 获取答案的完整链路追踪

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

📑 目录

1. 全流程概览 2. Step 1: 任务接入 (task_executor.py) 3. Step 2: 文档解析 (qa.py) 4. Step 3: Chunk 构建 (qa.py → beAdoc) 5. Step 4: LLM 增强 (关键词/问题/标签) 6. Step 5: 向量嵌入 (embedding) 7. Step 6: 存储入库 (insert_chunks → ES/Infinity) 8. Step 7: 查询构建 (query.py → FulltextQueryer) 9. Step 8: 多阶段检索 (search.py → Dealer) 10. Step 9: 重排序 & 引用标注 11. Step 10: LLM 上下文格式化 12. 完整数据流总结

1. 全流程概览

🔄 QA 文档端到端处理流水线

flowchart TB
    A["📥 文件上传 (Excel/CSV/PDF/MD/DOCX)"] --> B["Step 1: 任务接入
task_executor.collect()
Redis Stream 消费"] B --> C["Step 2: 文档解析
qa.chunk() → Excel/Pdf/Docx"] C --> D["Step 3: Chunk构建
beAdoc(Q,A) → content_with_weight + content_ltks"] D --> E["Step 4: LLM增强
keyword_extraction + question_proposal + content_tagging"] E --> F["Step 5: 向量嵌入
embedding() → q_X_vec"] F --> G["Step 6: 存储入库
insert_chunks() → ES/Infinity/OB"] G --> H["📊 文档存储
content_with_weight: Q+A
content_ltks: Q only
q_X_vec: Q+A embedding"] I["🔍 用户查询"] --> J["Step 7: 查询构建
FulltextQueryer.question()
field boosts: question_tks^20, content_ltks^2"] J --> K["Step 8: 多阶段检索
Dealer.retrieval()
全文+向量+融合"] K --> L["Step 9: 重排序
rerank_by_model() + insert_citations()"] L --> M["Step 10: LLM上下文
kb_prompt() 格式化 Q+A 送入 LLM"] M --> N["📋 带引用的答案"] H --> K style A fill:#dbeafe,stroke:#2563eb style H fill:#dcfce7,stroke:#16a34a style I fill:#f3e8ff,stroke:#9333ea style N fill:#fee2e2,stroke:#dc2626

2. Step 1: 任务接入 task_executor.py

🔌 入口:FACTORY 注册 + Redis Stream 消费

QA 解析器通过 FACTORY 字典注册为 ParserType.QA 的处理器:

# task_executor.py:104-121
FACTORY = {
    ParserType.QA.value: qa,   # ← QA 解析器注册
    ...
}

# task_executor.py:265-316  build_chunks()
chunker = FACTORY[task["parser_id"].lower()]  # 获取 qa 模块
binary = await get_storage_binary(bucket, name)  # 从 MinIO 拉取文件
cks = await thread_pool_exec(
    chunker.chunk,                              # 调用 qa.chunk()
    task["name"], binary=binary,
    from_page=task["from_page"], to_page=task["to_page"],
    lang=task_language, callback=progress_callback,
    kb_id=task["kb_id"], parser_config=parser_config_for_chunk,
    tenant_id=task["tenant_id"],
)

关键点:QA 任务与其他 parser 类型走同一入口,通过 parser_id="qa" 分发到 rag.app.qa.chunk()

3. Step 2: 文档解析 qa.py → chunk()

📄 chunk() 按文件扩展名分发到 6 种解析器

格式解析类/方法核心逻辑
.xlsxExcel.__call__()openpyxl 逐行读取,第1列=Q,第2列=A
.csvcsv.reader自动检测分隔符 (\t 或 ,),第1列=Q,第2列=A
.txtget_text() + split自动检测分隔符,第1列=Q,第2列=A
.pdfPdf.__call__()OCR → 布局分析 → Q-bullet 检测 → Q&A 配对
.mdmdQuestionLevel()Markdown 标题层级检测 Q,正文=A
.docxDocx.__call__()段落样式/大纲级别检测 Q,正文=A

🔑 PDF Q&A 检测核心:has_qbullet()

PDF 解析是最复杂的路径,使用 qbullets_category() 分析全文 bullets 模式,然后用 has_qbullet() 逐行判定是否为 Q 的起始行:

# qa.py:109-174
q_bull, reg = qbullets_category(sections)  # 分析全文的 Q bullet 正则模式
for box in self.boxes:
    has_bull, index = has_qbullet(reg, box, last_box, last_index, last_bull, bull_x0_list)
    if not has_bull:  # 不是 Q 的起始 → 追加到当前 A
        last_a = f'{last_a}{sum_section}'
    else:             # 是 Q 的起始 → 保存上一对 Q&A,开始新的 Q
        if last_q:
            qai_list.append((last_q, last_a, image, poss))
        last_q = has_bull.group()       # Q 文本
        last_a = section[end:]           # A 文本(Q bullet 之后的部分)

4. Step 3: Chunk 构建 qa.py → beAdoc()

🏗️ beAdoc() — Q&A chunk 的核心构建逻辑

⚡ 关键设计:content_with_weight = Q+A,content_ltks = 仅 Q
# qa.py:291-300
def beAdoc(d, q, a, eng, row_num=-1):
    qprefix = "Question: " if eng else "问题:"
    aprefix = "Answer: " if eng else "回答:"
    d["content_with_weight"] = "\t".join(
        [qprefix + rmPrefix(q), aprefix + rmPrefix(a)])  # ← Q+A 拼接(向量嵌入 + LLM 上下文)
    d["content_ltks"] = rag_tokenizer.tokenize(q)          # ← 仅 Q 分词(全文检索)
    d["content_sm_ltks"] = rag_tokenizer.fine_grained_tokenize(d["content_ltks"])  # ← 仅 Q 细粒度分词
    if row_num >= 0:
        d["top_int"] = [row_num]
    return d

字段语义表:

字段内容来源用途检索阶段
content_with_weight"问题:xxx\t回答:yyy"① 向量嵌入的原文 ② 返回给 LLM 的上下文向量检索 + LLM 上下文
content_ltks仅 Q 分词全文倒排索引的关键词全文检索
content_sm_ltks仅 Q 细粒度分词中文子词级匹配全文检索(子词)
imagePDF 中 Q&A 区域的截图多模态检索向量检索(如有 VLM 描述)
top_int行号 (Excel/CSV)排序/位置参考排序

5. Step 4: LLM 增强 task_executor.py

🤖 三个 LLM 后处理阶段(可选,由 parser_config 控制)

阶段触发条件函数生成字段对 QA 的影响
关键词提取auto_keywords > 0keyword_extraction()important_kwd, important_tks基于 Q+A 全文提取关键词,增强全文检索召回
问题生成auto_questions > 0question_proposal()question_kwd, question_tks基于 Q+A 生成更多等价问题,嵌入时优先使用 question_kwd
内容标签tag_kb_ids 非空content_tagging()tag_kwd自动分类标签,支持标签过滤检索
# task_executor.py:637-648 — embedding 函数中的关键逻辑
async def embedding(docs, mdl, parser_config=None, callback=None):
    for d in docs:
        tts.append(d.get("docnm_kwd", "Title"))
        c = "\n".join(d.get("question_kwd", []))       # ← 优先使用 LLM 生成的问题
        if not c:
            c = d["content_with_weight"]                # ← 降级使用 Q+A 原文
        cnts.append(c)
⚠️ 对 QA 的重要影响:如果开启了 auto_questions,embedding 时会用 LLM 生成的等价问题替代 Q+A 原文来做向量嵌入。这意味着向量检索的语义空间从"Q+A 文本"变成了"LLM 改写的问题集",可能提升对多样化自然语言 query 的召回率。

6. Step 5: 向量嵌入 task_executor.py → embedding()

📐 双向量加权融合

# task_executor.py:637-686
async def embedding(docs, mdl, parser_config=None, callback=None):
    # 1. 准备文本
    for d in docs:
        tts.append(d.get("docnm_kwd", "Title"))         # 标题文本
        c = "\n".join(d.get("question_kwd", []))         # LLM生成的问题(优先)
        if not c:
            c = d["content_with_weight"]                  # 降级:Q+A 原文
        cnts.append(c)

    # 2. 编码标题向量
    vts, c = await thread_pool_exec(mdl.encode, tts[0:1])
    tts = np.tile(vts[0], (len(cnts), 1))               # 所有 chunk 共用标题向量

    # 3. 批量编码内容向量
    for batch in cnts_batches:
        vts, c = await thread_pool_exec(batch_encode, batch)
        cnts_batches.append(vts)

    # 4. 加权融合
    filename_embd_weight = parser_config.get("filename_embd_weight", 0.1)
    title_w = float(filename_embd_weight)                # 默认 0.1
    vects = title_w * tts + (1 - title_w) * cnts        # 标题向量×0.1 + 内容向量×0.9

    # 5. 写入 chunk
    for i, d in enumerate(docs):
        v = vects[i].tolist()
        d["q_%d_vec" % len(v)] = v                       # 例: q_1024_vec

对 QA 的关键影响:

7. Step 6: 存储入库 insert_chunks()

💾 批量写入 ES / Infinity / OceanBase

# task_executor.py:1174-1272
async def insert_chunks(task_id, task_tenant_id, task_dataset_id, chunks, progress_callback):
    # 1. 处理母 chunk(mom/mom_with_weight)
    for ck in chunks:
        mom = ck.get("mom") or ck.get("mom_with_weight") or ""
        if mom:
            # 创建 mom chunk(仅保留必要字段,available_int=0 不可被检索)
            mothers.append(mom_ck)

    # 2. 先插入母 chunk
    for b in range(0, len(mothers), settings.DOC_BULK_SIZE):
        await thread_pool_exec(settings.docStoreConn.insert, ...)

    # 3. 批量插入 QA chunks
    for b in range(0, len(chunks), settings.DOC_BULK_SIZE):
        doc_store_result = await thread_pool_exec(
            settings.docStoreConn.insert,
            chunks[b:b + settings.DOC_BULK_SIZE],
            search.index_name(task_tenant_id),
            task_dataset_id,
        )
        # 取消时回滚已插入的 chunk

QA chunk 在 ES/Infinity 中的最终字段:

字段值示例索引方式
idxxhash(Q+A+doc_id)主键
content_with_weight"问题:如何退款\t回答:请登录..."不直接索引(用于 embedding 原文)
content_ltks"如何 退款"全文倒排索引
content_sm_ltks"如 何 退 款"细粒度倒排索引
important_kwd["退款", "订单", "流程"]关键词字段 (boost 30x)
important_tksLLM提取的关键词分词关键词分词索引 (boost 20x)
question_kwd["怎么申请退款", "退款流程"]LLM生成的等价问题
question_tksLLM生成问题的分词问题分词索引 (boost 20x)
q_1024_vec[0.12, -0.34, ...]向量索引 (cosine)
doc_id"abc123"过滤字段
kb_id["kb_001"]过滤字段

8. Step 7: 查询构建 query.py → FulltextQueryer.question()

🔍 用户 query → 结构化全文查询表达式

# query.py:28-41 — 字段权重配置
class FulltextQueryer(QueryBase):
    def __init__(self):
        self.query_fields = [
            "title_tks^10",          # 标题分词 boost=10
            "title_sm_tks^5",        # 标题细粒度 boost=5
            "important_kwd^30",      # LLM 关键词 boost=30 ← 最高权重!
            "important_tks^20",      # LLM 关键词分词 boost=20
            "question_tks^20",       # LLM 生成问题分词 boost=20
            "content_ltks^2",        # Q-only 内容分词 boost=2
            "content_sm_ltks",       # Q-only 细粒度 boost=1(默认)
        ]

    def question(self, txt, tbl="qa", min_match: float = 0.6):
        # 1. 文本清洗:繁简转换、全半角转换、特殊字符过滤
        txt = rag_tokenizer.tradi2simp(rag_tokenizer.strQ2B(txt.lower()))
        
        # 2. 英文路径:词权重 + 同义词扩展 + bigram
        if not self.is_chinese(txt):
            tks_w = self.tw.weights(tks, preprocess=False)  # 词权重
            syns = [self.syn.lookup(tk) for tk, w in tks_w]  # 同义词
            # 构建 "(tk^weight synonym_expansion)" 查询表达式
        
        # 3. 中文路径:细粒度分词 + 同义词 + 子词 n-gram
        # 构建匹配表达式
        
        # 4. 返回 MatchTextExpr → 用于 ES/Infinity 全文检索
🔑 对 QA 的关键影响:
  • content_ltks^2 — QA chunk 的 content_ltks 只包含 Q 的分词,boost=2 相对较低
  • important_kwd^30 — 如果开启了 auto_keywords,LLM 提取的关键词获得最高权重
  • question_tks^20 — 如果开启了 auto_questions,LLM 生成的问题获得高权重
  • 这意味着:全文检索阶段,LLM 增强字段的权重远超原始 Q-only 字段

9. Step 8: 多阶段检索 search.py → Dealer.retrieval()

🎯 retrieval() — 6 阶段检索引擎

检索流水线时序

sequenceDiagram
    participant U as 用户查询
    participant D as Dealer.retrieval()
    participant Q as FulltextQueryer
    participant DS as DocStore (ES/Infinity)
    participant E as EmbeddingModel
    participant R as RerankModel

    U->>D: retrieval(question, embd_mdl, ...)
    D->>Q: question(txt) → MatchTextExpr
    D->>E: encode_queries(txt) → MatchDenseExpr
    D->>DS: search(text_expr, dense_expr, fusion)
    DS-->>D: 初检结果 (top_k × 2)
    D->>D: _prune_deleted_chunks()
    alt rerank_by_model
        D->>R: similarity(query, chunks)
        R-->>D: 精排分数
    else 本地重排
        D->>D: rerank() / rerank_with_knn()
    end
    D->>D: 按 similarity 阈值过滤
    D->>D: 分页 + group_docs 聚合
    D-->>U: SearchResult (ids, field, highlight, group_docs)
    

QA 检索的混合匹配逻辑:

  1. 全文检索:用户 query 分词 → 匹配 content_ltks(Q-only 分词)、important_tksquestion_tks 等字段
  2. 向量检索:用户 query → embedding → 与 q_X_vec(Q+A 向量)做 cosine 相似度
  3. 融合:全文分数 + 向量分数按权重合并
  4. 精排:可选 Cross-Encoder 模型对 Top-N 做精确重排序
  5. 过滤:相似度低于 threshold 的 chunks 被丢弃
  6. 聚合:同一文档的 chunks 合并为 group_docs

10. Step 9: 重排序 & 引用标注

📎 rerank_by_model() + insert_citations()

重排序(可选):

# search.py — 使用 Cross-Encoder 模型精排
def rerank_by_model(self, rerank_mdl, sres, query, tkweight, vtweight, cfield, rank_feature):
    # 对初检结果的每个 chunk 文本,用 Cross-Encoder 模型重新打分
    # 结合 token 相似度和向量相似度
    ...
    return sres  # 按新分数重排

引用标注:

# search.py:242-332  insert_citations()
# 1. 将 LLM 回答按句子切分
# 2. 对每个句子编码为向量
# 3. 计算句子与检索到的 chunks 的混合相似度
# 4. 对相似度 > 0.63 的句子标注 [ID:n] 引用
# 5. 阈值逐步降低(0.63 → 0.50 → 0.40...)直到找到引用
for i, a in enumerate(pieces_):
    sim, tksim, vtsim = self.qryr.hybrid_similarity(
        ans_v[i], chunk_v, a_tks, chunks_tks, tkweight, vtweight)
    if mx >= thr:
        cites[idx[i]] = [str(ii) for ii in range(len(chunk_v)) if sim[ii] > mx]
# 输出: "退款请登录账户...[ID:3] 进入订单页面...[ID:7]"
对 QA 的意义:对于 FAQ 场景,LLM 可能会引用多个 Q&A chunks 来回答。引用标注机制能精确定位答案来自哪些 Q&A 对。

11. Step 10: LLM 上下文格式化 generator.py → kb_prompt()

📝 检索结果 → LLM Prompt 的格式化

# prompts/generator.py  kb_prompt()
def kb_prompt(kbinfos, max_tokens, hash_id):
    # 将检索到的 chunks 格式化为结构化 Prompt:
    # [ID]: 1
    # [Title]: 常见问题.xlsx
    # [Content]: 问题:如何退款?\t回答:请登录账户,进入订单页面...
    # [Metadata]: {"tag": "退款"}
    #
    # [ID]: 2
    # ...
    # Token 预算控制:超过 max_tokens 时截断

LLM 收到的上下文包含完整的 Q+A,可以直接基于标准答案回复,而非重新生成。

12. 完整数据流总结

🧬 QA Chunk 字段在全流程中的演变

flowchart LR
    subgraph Parse["解析阶段"]
        P1["Excel/CSV: Q列+A列"]
        P2["PDF: Q-bullet检测+A区域"]
        P3["MD/DOCX: 标题层级检测"]
    end

    subgraph Build["Chunk构建 (beAdoc)"]
        B1["content_with_weight = 问题:Q \\t 回答:A"]
        B2["content_ltks = tokenize(Q)"]
        B3["content_sm_ltks = fine_grained_tokenize(Q)"]
    end

    subgraph Enhance["LLM增强"]
        E1["important_kwd ← keyword_extraction(Q+A)"]
        E2["question_kwd ← question_proposal(Q+A)"]
        E3["tag_kwd ← content_tagging(Q+A)"]
    end

    subgraph Embed["向量嵌入"]
        EM1["优先: encode(question_kwd)"]
        EM2["降级: encode(content_with_weight)"]
        EM3["q_X_vec = title_vec×0.1 + content_vec×0.9"]
    end

    subgraph Search["检索阶段"]
        S1["全文: content_ltks(boost2) + important_kwd(boost30) + question_tks(boost20)"]
        S2["向量: q_X_vec cosine similarity"]
        S3["融合: text_score + vector_score"]
    end

    P1 --> B1
    P2 --> B1
    P3 --> B1
    B1 --> E1
    B1 --> E2
    B1 --> E3
    B2 --> S1
    E1 --> EM1
    E1 --> S1
    E2 --> EM1
    E2 --> S1
    EM1 --> EM3
    EM2 --> EM3
    EM3 --> S2
    S1 --> S3
    S2 --> S3

    style B1 fill:#dbeafe,stroke:#2563eb
    style B2 fill:#ffedd5,stroke:#ea580c
    style EM3 fill:#dcfce7,stroke:#16a34a
    style S3 fill:#fee2e2,stroke:#dc2626
  

📊 关键设计决策与影响分析

设计决策实现方式对检索的影响
Q+A 一起存content_with_weight = Q+A✅ 向量检索利用 A 的语义丰富性
✅ LLM 直接获得标准答案
仅 Q 做全文分词content_ltks = tokenize(Q)✅ 全文检索精确匹配在 Q 上
✅ 避免 A 文本引入噪声
LLM 关键词高权重important_kwd^30✅ 弥补 Q 文本短导致的匹配不足
⚠️ 依赖 LLM 提取质量
LLM 问题生成优先嵌入embedding() 优先用 question_kwd✅ 嵌入空间对齐用户自然语言 query
⚠️ 多消耗 LLM token
双向量加权融合title_w=0.1, content_w=0.9✅ 轻度约束同文档 chunks 的向量方向
✅ 内容语义主导
Cross-Encoder 精排(可选)rerank_by_model()✅ 大幅提升 Top-N 精度
⚠️ 增加检索延迟