RagFlow 0.25.6 全流程源码分析

前端 → 后端 → 数据存储完整链路追踪,含关键源码与行号
源码根目录: /home/yzq/tmp/ragflow/auth/ragflow-0.25.6

整体架构与组件

RagFlow 是一个基于深度文档理解的 RAG(检索增强生成)引擎。其核心数据流可分为两条主线:文件入库(上传→对象存储→解析→向量库)与检索召回(查询→混合检索→重排→返回)。

组件职责
前端React + Vite + React Query表单采集、FormData 上传、检索参数提交
APIQuart (async Flask) api/apps/restful_apis/鉴权、校验、参数解析、编排 Service
业务api/db/services/文件/文档/任务的业务逻辑
对象存储MinIO (默认) / S3 / OSS / GCS / Azure存储原始文件二进制、缩略图、chunk 图片
关系库MySQL (Peewee ORM)document / file / file2document / task
消息队列Redis (consumer group)异步解析任务分发
向量库Elasticsearch (默认) / Infinity / OpenSearch / OceanBasechunk 全文 + 稠密向量混合索引 ragflow_{tenant_id}
版本差异提示:此版本路由层位于 api/apps/restful_apis/(非社区标准版的 document_app.py)。rag/utils/storage_factory.py 在本版本为空文件,真正的 StorageFactory 定义在 common/settings.py

架构图 (组件拓扑)

两条主线共享同一套存储与服务组件:实线为同步调用,虚线为异步消息流。

flowchart TB subgraph FE["前端 React + Vite"] UP["上传组件
use-document-request"] SE["检索表单
testing-form"] end subgraph API["API 层 Quart · api/apps/restful_apis"] DOC["document_api
upload_document"] CHK["chunk_api
retrieval_test"] end subgraph SVC["业务层 api/db/services"] FS["FileService"] DS["DocumentService"] TS["TaskService
queue_tasks"] RT["Dealer.retrieval
rag/nlp/search"] end subgraph STORE["存储与中间件"] MINIO[("MinIO
bucket=kb_id")] MYSQL[("MySQL
document/file/task")] REDIS[["Redis 队列"]] VDB[("ES / Infinity
ragflow_tenant_id")] end EXE["task_executor
build_chunks + embedding"] EMB["Embedding 模型"] RRK["Rerank 模型"] UP -->|"FormData POST"| DOC SE -->|"POST /datasets/search"| CHK DOC --> FS FS -->|"put 原始文件"| MINIO FS -->|"insert 元数据"| MYSQL DOC -.触发解析.-> DS DS --> TS TS -->|"切分任务"| MYSQL TS -.推入队列.-> REDIS REDIS -.消费.-> EXE EXE -->|"取回原文件"| MINIO EXE --> EMB EXE -->|"insert chunks"| VDB CHK --> RT RT --> EMB RT -->|"混合检索"| VDB RT --> RRK classDef fe fill:#16321f,stroke:#7ee787,color:#d6deeb; classDef api fill:#102a44,stroke:#4ea1ff,color:#d6deeb; classDef store fill:#3a2a18,stroke:#ffab70,color:#d6deeb; class UP,SE fe; class DOC,CHK,FS,DS,TS,RT,EXE api; class MINIO,MYSQL,REDIS,VDB store;

一、文件上传到存储全流程

前端 FormData ──POST /api/v1/datasets/{id}/documents──▶ document_api.upload_document │ ├─▶ FileService.upload_document │ ├─▶ MinIO.put(bucket=kb_id, key=location) [原始文件二进制] │ └─▶ MySQL: document + file + file2document [元数据] │ └─(触发解析)─▶ DocumentService.run ─▶ queue_tasks ─▶ Redis 队列 │ ▼ task_executor 消费 ─▶ build_chunks(切片) ─▶ embedding(向量化) │ ▼ docStoreConn.insert ─▶ ES/Infinity 索引 ragflow_{tenant_id}

▸ 上传时序图

sequenceDiagram autonumber participant U as 用户/浏览器 participant FE as 前端 React
(use-document-request) participant API as document_api.py
upload_document participant FS as FileService
upload_document participant OSS as MinIO
(bucket=kb_id) participant DB as MySQL
(document/file/file2document) participant DS as DocumentService.run
+ queue_tasks participant MQ as Redis 队列 participant TE as task_executor participant VDB as ES/Infinity
ragflow_{tenant_id} U->>FE: 选择文件 FE->>FE: 构造 FormData (file[], parser_config) FE->>API: POST /api/v1/datasets/{id}/documents API->>API: 鉴权 / KB 校验 / 团队权限 / 文件名长度 API->>FS: thread_pool_exec(upload_document) FS->>OSS: put(kb_id, location, blob) + 缩略图 OSS-->>FS: ok (重名追加 "_") FS->>DB: DocumentService.insert + add_file_from_kb DB-->>FS: 原子递增 KB 文档数 FS-->>API: (doc, blob) API-->>FE: code=0 + 文档元数据 alt 勾选 parseOnCreation FE->>API: runDocumentByIds(run=1) API->>DS: DocumentService.run DS->>OSS: 取 PDF 算页数 (切分依据) DS->>MQ: queue_product(切分后的 task) MQ-->>TE: consumer group 拉取 TE->>OSS: get(bucket, name) 取回原文件 TE->>TE: build_chunks 切片 + embedding 向量化 TE->>VDB: docStoreConn.insert(chunks) TE->>DB: increment_chunk_num / 更新 progress end
1前端上传组件与 API 调用 前端
用户选择文件,构造 FormData,axios POST 到后端;成功后可选触发解析。

URL 定义

web/src/utils/api.ts:2, 132-133
const restAPIv1 = `/api/v1`;
documentUpload: (datasetId: string) =>
  `${restAPIv1}/datasets/${datasetId}/documents`,

构造 FormData 并发请求

web/src/hooks/use-document-request.ts:82-91
const formData = new FormData();
fileList.forEach((file) => {
  formData.append('file', file);          // 字段名固定 "file",支持多文件
});
if (parserConfig) {
  formData.append('parser_config', JSON.stringify(parserConfig));
}
const ret = await uploadDocument(id, formData);

底层请求(裸 axios 以手动设置 multipart 头)

web/src/services/knowledge-service.ts:321-329
export const uploadDocument = async (datasetId, formData) => {
  const url = api.documentUpload(datasetId);
  const response = await axios.post(url, formData, {
    headers: { [Authorization]: getAuthorization() },
  });
  return response.data;
};

上传成功(code===0)且勾选 parseOnCreation 时,调用 runDocumentByIds({documentIds, run:1}) 触发解析。请求体为 multipart:file(可多个)+ 可选 parser_config;query 可带 ?type=web|empty|local

2后端 API 接收文件 API
校验 KB 存在与团队权限、文件名长度,按上传类型分流,解析 parser_config 白名单覆盖项。
api/apps/restful_apis/document_api.py:366-458 — 主路由
@manager.route("/datasets/<dataset_id>/documents", methods=["POST"])
@login_required
@add_tenant_id_to_kwargs
async def upload_document(dataset_id, tenant_id):
    upload_type = (request.args.get("type") or "local").lower()
    e, kb = KnowledgebaseService.get_by_id(dataset_id)
    if not e: ...                                  # KB 不存在
    if not check_kb_team_permission(kb, tenant_id): ... # 权限校验
    if upload_type == "web":   return await _upload_web_document(...)
    if upload_type == "empty": return await _upload_empty_document(...)
    return await _upload_local_documents(kb, tenant_id)
api/apps/restful_apis/document_api.py:571-607 — 本地文件接收
async def _upload_local_documents(kb, tenant_id):
    files = await request.files
    if "file" not in files:
        return get_error_data_result(message="No file part!")
    file_objs = files.getlist("file")
    for file_obj in file_objs:
        if len(file_obj.filename.encode("utf-8")) > FILE_NAME_LEN_LIMIT: ...
    allowed_keys = {"table_column_mode", "table_column_roles"}  # 白名单
    err, files = await thread_pool_exec(          # 同步阻塞放线程池
        FileService.upload_document, kb, file_objs, tenant_id,
        parent_path=form.get("parent_path"),
        parser_config_override=parser_config_override,
    )

此处只接收落库,不触发解析。解析由独立 parse 接口触发(见阶段 5)。

3文件存储到对象存储 (MinIO) 存储
写入对象存储;bucket = 知识库 ID,key = 文件名,重名追加 "_",并生成缩略图。
common/settings.py:181-194 — 存储工厂(真正的工厂)
class StorageFactory:
    storage_mapping = {
        Storage.MINIO: RAGFlowMinio,  Storage.AWS_S3: RAGFlowS3,
        Storage.OSS: RAGFlowOSS, Storage.OPENDAL: OpenDALStorage,
        Storage.GCS: RAGFlowGCS, ...
    }
    @classmethod
    def create(cls, storage): return cls.storage_mapping[storage]()
api/db/services/file_service.py:512-531 — 核心存储调用("存到哪")
filename = duplicate_name(DocumentService.query, name=file.filename, kb_id=kb.id)
location = filename if not safe_parent_path else f"{safe_parent_path}/{filename}"
while settings.STORAGE_IMPL.obj_exist(kb.id, location):   # bucket = kb.id
    location += "_"                                       # 重名追加下划线
blob = file.read()
if filetype == FileType.PDF.value:
    blob = read_potential_broken_pdf(blob)
settings.STORAGE_IMPL.put(kb.id, location, blob)          # ← 写入对象存储
img = thumbnail_img(filename, blob)
if img is not None:
    thumbnail_location = f"thumbnail_{doc_id}.png"
    settings.STORAGE_IMPL.put(kb.id, thumbnail_location, img)
rag/utils/minio_conn.py:144-161 — MinIO put 实现
@use_default_bucket
@use_prefix_path
def put(self, bucket, fnm, binary, tenant_id=None):
    for _ in range(3):                                   # 失败重试 3 次
        if not self.bucket and not self.conn.bucket_exists(bucket):
            self.conn.make_bucket(bucket)                # 多桶模式按 kb_id 建桶
        r = self.conn.put_object(bucket, fnm, BytesIO(binary), len(binary))
        return r
bucket 命名规则: bucket = kb.id(知识库 ID),key = location。两个装饰器实现单桶/多桶部署:单桶模式把 kb_id 拼进 key 形成 <prefix>/<kb_id>/<fnm>
4元数据写入数据库 (MySQL) 存储
写入 document / file / file2document 三表,并原子递增 KB 文档计数。
api/db/db_models.py:904-931 — Document 模型核心字段
class Document(DataBaseModel):
    id = CharField(max_length=32, primary_key=True)
    kb_id = CharField(...)                    # 所属 KB = MinIO bucket
    parser_id = CharField(...)                # naive/table/paper/picture
    parser_config = JSONField(...)
    type / name / location / size / suffix    # location = MinIO object key
    token_num / chunk_num / progress / progress_msg
    content_hash = CharField(...)             # xxhash128 变更检测
    run = CharField(default="0")              # 0未跑/1运行/2取消
api/db/services/file_service.py:533-553 — 组装并插入
doc = {
    "id": doc_id, "kb_id": kb.id,
    "parser_id": self.get_parser(filetype, filename, kb.parser_id),
    "location": location, "size": len(blob),
    "thumbnail": thumbnail_location,
    "content_hash": incoming_fp or xxhash.xxh128(blob).hexdigest(),
}
DocumentService.insert(doc)                              # 写 document 表
FileService.add_file_from_kb(doc, kb_folder["id"], kb.tenant_id) # 写 file + file2document
api/db/services/file2document_service.py:82-96 — 地址反查 (下载/解析复用)
def get_storage_address(cls, doc_id=None, file_id=None):
    f2d = cls.get_by_document_id(doc_id)
    if f2d:
        file = File.get_by_id(f2d[0].file_id)
        if not file.source_type or file.source_type == FileSource.LOCAL:
            return file.parent_id, file.location
    e, doc = DocumentService.get_by_id(doc_id)
    return doc.kb_id, doc.location          # bucket=kb_id, key=location
5文档解析与切片入向量库 异步 存储
parse 接口标记 RUNNING → queue_tasks 切分(PDF 按页/Excel 按行)入 Redis → task_executor 消费切片向量化写入 ES/Infinity。
api/db/services/document_service.py:1044-1064 — DocumentService.run
if doc.get("pipeline_id", ""):
    queue_dataflow(tenant_id, flow_id=doc["pipeline_id"], task_id=get_uuid(), doc_id=doc["id"])
else:
    bucket, name = File2DocumentService.get_storage_address(doc_id=doc["id"])
    queue_tasks(doc, bucket, name, 0)
api/db/services/task_service.py:356-462 — queue_tasks 切分 + 入队
if doc["type"] == FileType.PDF.value:
    pages = PdfParser.total_page_number(doc["name"], file_bin)
    page_size = doc["parser_config"].get("task_page_size") or 12
    for ...: task["from_page"], task["to_page"] = ...   # 按页切任务
elif doc["parser_id"] == "table":
    rn = RAGFlowExcelParser.row_number(...)                  # Excel 按 3000 行切
bulk_insert_into_db(Task, parse_task_array, True)            # 写 task 表
for unfinished_task in unfinished_task_array:
    REDIS_CONN.queue_product(settings.get_svr_queue_name(priority), message=unfinished_task)
rag/svr/task_executor.py:265-323 — build_chunks 下载并切片
chunker = FACTORY[task["parser_id"].lower()]
bucket, name = File2DocumentService.get_storage_address(doc_id=task["doc_id"])
binary = await get_storage_binary(bucket, name)            # 从 MinIO 取回原文件
cks = await thread_pool_exec(chunker.chunk, task["name"], binary=binary,
        from_page=task["from_page"], to_page=task["to_page"],
        kb_id=task["kb_id"], parser_config=parser_config_for_chunk)
rag/svr/task_executor.py:631-634 / 1447-1505 — 向量化并写入向量库
def init_kb(row, vector_size):
    idxnm = search.index_name(row["tenant_id"])            # ragflow_{tenant_id}
    return settings.docStoreConn.create_idx(idxnm, row.get("kb_id",""), vector_size, parser_id)
# do_handle_task 主流程:
chunks = await build_chunks(task, progress_callback)
token_count, vector_size = await embedding(chunks, embedding_model, ...)
await thread_pool_exec(settings.docStoreConn.insert,
        chunks[b:b+settings.DOC_BULK_SIZE],
        search.index_name(task_tenant_id), task_dataset_id)   # ← 写入 ES/Infinity
DocumentService.increment_chunk_num(task_doc_id, task_dataset_id, token_count, chunk_count, 0)

核心标识符规则 & 设计点

[testing-form.tsx] question/threshold/weight/top_k │ useWatch → setValues [use-knowledge-request.ts] +page/size/doc_ids/highlight → kbService.retrievalTest [knowledge-service.ts] dataset_id → dataset_ids[] POST /datasets/search ▼ [chunk_api.py retrieval_test] 鉴权/校验/装载 embd+rerank/可选 query 预处理 │ settings.retriever.retrieval(...) ▼ [search.py Dealer.retrieval] 算 RERANK_LIMIT/分页 → search() ├─ query.py question 全文表达式 MatchTextExpr + keywords ├─ get_vector encode_queries → MatchDenseExpr(q_dim_vec) └─ FusionExpr weighted_sum 0.05,0.95 ▼ [es_conn.py search] query_string + knn + rank_feature → SearchResult ├─ 重排分支: rerank_mdl / Infinity(atan) / OceanBase / ES(二次KNN) │ sim = tkweight*term + vtweight*vector + rank_feature ├─ argsort + similarity_threshold 过滤 + 分页切片 └─ chunks{similarity,vector_similarity,term_similarity,highlight} + doc_aggs ▼ [chunk_api.py] 字段名映射 → get_result → 前端渲染
sequenceDiagram autonumber participant U as 用户 participant FE as 前端
testing-form / hook participant SVC as service
knowledge-service.ts participant API as chunk_api.py
retrieval_test participant EMB as Embedding 模型 participant QRY as query.py
FulltextQueryer participant DEAL as search.py
Dealer participant DS as 向量库
ES / Infinity / OceanBase participant RR as Rerank 模型(可选) U->>FE: 输入 question + 阈值/权重/top_k FE->>FE: useWatch 实时回写表单值 FE->>SVC: retrievalTest(params) SVC->>SVC: dataset_id → dataset_ids[] SVC->>API: POST /datasets/search API->>API: 鉴权 + 校验 embedding 模型一致 API->>EMB: 装载 embedding(必) + rerank(可选) API->>DEAL: retriever.retrieval(...) DEAL->>QRY: question() 改写为带权布尔表达式 QRY-->>DEAL: MatchTextExpr + keywords DEAL->>EMB: encode_queries(question) EMB-->>DEAL: 稠密向量 → MatchDenseExpr(q_dim_vec) DEAL->>DS: search(全文 + KNN + FusionExpr 0.05,0.95) DS-->>DEAL: SearchResult(命中候选) alt 空结果 DEAL->>DS: 降级重试(min_match↓, similarity↑) end alt 有 rerank 模型 DEAL->>RR: similarity(query, docs) RR-->>DEAL: 向量相似度分 else ES 默认 DEAL->>DS: 二次 KNN-only 取干净 cosine (_knn_scores) DS-->>DEAL: cosine 分 end DEAL->>DEAL: sim = tk*term + vt*vector + rank_feature
阈值过滤 + 排序 + 分页 DEAL-->>API: chunks{3类相似度,highlight} + doc_aggs API->>API: 去 vector 字段 + 字段名映射 API-->>FE: get_result(data=ranks) FE-->>U: 渲染 chunk 列表 / 高亮 / 文档聚合
1前端检索入口 前端
表单采集检索参数,useWatch 实时回写,hook 拼装分页+highlight,service 归一 dataset_ids。
web/src/pages/dataset/testing/testing-form.tsx:57-86 — 表单 schema
const formSchema = z.object({
  question: z.string().min(1, ...),
  ...similarityThresholdSchema,         // similarity_threshold
  ...vectorSimilarityWeightSchema,      // vector_similarity_weight
  ...topKSchema,                        // top_k
  use_kg: z.boolean().optional(),
  dataset_ids: z.array(z.string()).optional(),
});
const values = useWatch({ control: form.control });
useEffect(() => { setValues(values); }, [setValues, values]);
web/src/hooks/use-knowledge-request.ts:62-93 — useTestRetrieval
const queryParams = {
    ...values,
    kb_id: values?.kb_id || knowledgeBaseId,
    page, size: pageSize,            // 默认 page=1, pageSize=10
    doc_ids: filterValue.doc_ids,    // 文档级过滤
    highlight: true,
};
mutationFn: async (params) => {
    const { data } = await kbService.retrievalTest(params);
    return { ...data?.data, isRuned: true };
}
web/src/services/knowledge-service.ts:150-165 + utils/api.ts:112
retrievalTest: async (params) => {
  const datasetId = params.dataset_id || params.kb_id || params.knowledge_id;
  const datasetIds = Array.isArray(datasetId) ? datasetId : [datasetId];
  return request.post(api.retrievalTest, {
    data: { ...rest, dataset_ids: datasetIds },   // 归一成数组
  });
},
// api.ts:  retrievalTest: `${restAPIv1}/datasets/search`
2后端 API 接收 API
鉴权 → 校验 embedding 模型一致 → 装载 embd/rerank → 可选 query 预处理 → 调核心检索 → 字段重映射。
api/apps/restful_apis/chunk_api.py:228-347 — retrieval_test
page = int(req.get("page", 1))
size = int(req.get("page_size", 30))
similarity_threshold = float(req.get("similarity_threshold", 0.2))
vector_similarity_weight = float(req.get("vector_similarity_weight", 0.3))
top = int(req.get("top_k", 1024))
# 校验所有 dataset 的 embedding 模型必须一致
embd_nms = list(set([... for kb in kbs]))
if len(embd_nms) != 1:
    return get_result(message="Datasets use different embedding models.", ...)
# 核心调用
ranks = await settings.retriever.retrieval(
    question, embd_mdl, tenant_ids, kb_ids, page, size, similarity_threshold,
    vector_similarity_weight, top, doc_ids, rerank_mdl=rerank_mdl,
    highlight=highlight, rank_feature=label_question(question, kbs))
chunk_api.py:334-343 — 返回前字段映射(内部名→前端契约)
key_mapping = {
    "chunk_id": "id", "content_with_weight": "content",
    "doc_id": "document_id", "important_kwd": "important_keywords",
    "docnm_kwd": "document_keyword", "kb_id": "dataset_id",
}
ranks["chunks"] = [{key_mapping.get(k, k): v for k,v in c.items()} for c in ranks["chunks"]]

对话(chat)路径走同一个 retriever.retrievaldialog_service.py:714-729),参数来自 dialog 配置,结果经 kb_prompt 拼成 LLM 上下文。

3Query 理解与 Embedding 编码 核心
全文查询改写成带权布尔表达式;问句编码成稠密向量。
rag/nlp/query.py:32-95 — FulltextQueryer 字段加权 + 改写
self.query_fields = [
    "title_ks^10", "title_sm_tks^5",
    "important_kwd^30", "important_tks^20", "question_tks^20",
    "content_ltks^2", "content_sm_ltks",
]
# 分词→term weight 加权→同义词扩展(1/4权重)→相邻词组 bigram(*2)
syn = ["\"{}\"^{:.4f}".format(s, w / 4.) for s in syn ...]
return MatchTextExpr(self.query_fields, query, 100, {...}), keywords
rag/nlp/search.py:53-61 — get_vector 编码查询向量
async def get_vector(self, txt, emb_mdl, topk=10, similarity=0.1):
    qv, _ = await thread_pool_exec(emb_mdl.encode_queries, txt)
    embedding_data = [get_float(v) for v in qv]
    vector_column_name = f"q_{len(embedding_data)}_vec"   # 维度编进列名 q_1024_vec
    return MatchDenseExpr(vector_column_name, embedding_data, 'float', 'cosine', topk, ...)
4混合检索核心 向量库
三段式:全文 query_string + 稠密向量 knn + FusionExpr 加权融合;空结果降级重试。
rag/nlp/search.py:132-212 — Dealer.search 编排
matchDense = await self.get_vector(qst, emb_mdl, topk, req.get("similarity", 0.1))
fusionExpr = FusionExpr("weighted_sum", topk, {"weights": "0.05,0.95"})
matchExprs = [matchText, matchDense, fusionExpr]
res = await thread_pool_exec(self.dataStore.search, src, highlightFields, filters,
                             matchExprs, orderBy, offset, limit, idx_names, kb_ids, ...)
# 空结果降级重试: min_match 0.3→0.1, similarity→0.17 再查一次
rag/utils/es_conn.py:194-230 — ES 混合查询实现
weights = m.fusion_params["weights"]
vector_similarity_weight = get_float(weights.split(",")[1])   # 取 0.95
# 全文 query_string,按 1 - vector_weight 调 boost
bool_query.must.append(Q("query_string", fields=m.fields, query=m.matching_text,
                         minimum_should_match=minimum_should_match, boost=1))
bool_query.boost = 1.0 - vector_similarity_weight
# 向量 KNN
s = s.knn(m.vector_column_name, m.topn, m.topn * 2,
          query_vector=list(m.embedding_data), filter=bool_query.to_dict(), similarity=similarity)
# rank_feature: PageRank/标签特征加成
bool_query.should.append(Q("rank_feature", field=fld, linear={}, boost=sc))
Infinity 后端融合用 atan 归一化(每路预归一后再融合),上层无需本地 rerank;OceanBase 仍回传向量走本地 rerank。
5Rerank 重排 核心
按后端/模型分四条路径;最终分 = term×权重 + vector×权重 + rank_feature。
rag/nlp/search.py:586-669 — retrieval 编排重排路径
RERANK_LIMIT = max(30, ceil(64/page_size)*page_size)
if rerank_mdl and top > 0: RERANK_LIMIT = min(RERANK_LIMIT, top, 64)
if rerank_mdl and sres.total > 0:
    sim, tsim, vsim = self.rerank_by_model(rerank_mdl, sres, question, ...)  # 外部模型
elif settings.DOC_ENGINE_INFINITY:
    sim = [sres.field[id].get("_score", 0.0) for id in sres.ids]  # 已归一直接用
elif settings.DOC_ENGINE_OCEANBASE:
    sim, tsim, vsim = self.rerank(sres, question, ...)                # 本地 rerank
else:                                                              # ES: 二次 KNN
    knn_scores = await self._knn_scores(sres, idx_names, kb_ids)
    sim, tsim, vsim = self.rerank_with_knn(sres, question, knn_scores, ...)
rag/nlp/search.py:443-472 — rerank_with_knn 最终加权公式
tksim = self.qryr.token_similarity(keywords, ins_tw)              # term 相似度
vtsim = [kn_scores.get(cid, 0.0) for cid in sres.ids]          # 向量 cosine(二次KNN)
rank_fea = self._rank_feature_scores(rank_feature, sres)         # PageRank/标签
sim = tkweight * tksim + vtweight * vtsim + rank_fea
# 分词字段加权: content_ltks + title*2 + important_kwd*5 + question_tks*6
设计亮点:ES 路径主检索不回传 chunk 向量,改用「候选 id 过滤的二次 KNN-only 查询」(_knn_scores) 取干净 cosine,再本地与 term 相似度按用户权重加权。外部 reranker 时向量分由模型 rerank_mdl.similarity 给出。
6分页过滤与结果返回 API
按总分降序排序、阈值过滤、分页切片,组装三类相似度分 + 文档聚合。
rag/nlp/search.py:676-733 — 排序/过滤/分页/组装
sorted_idx = np.argsort(sim_np * -1)                       # 按总分降序
valid_idx = [i for i in sorted_idx if sim_np[i] >= post_threshold]  # 阈值过滤
page_idx = valid_idx[begin:end]                              # 当前页切片
d = {
    "chunk_id": id, "content_with_weight": ...,
    "similarity": float(sim_np[i]),            # 综合分
    "vector_similarity": float(vsim[i]),       # 向量分
    "term_similarity": float(tsim[i]),         # 关键词分
}
if highlight and sres.highlight:
    d["highlight"] = remove_redundant_spaces(sres.highlight[id])

文档聚合 doc_aggs(按文档统计命中 chunk 数)供前端「按文档分组」展示。最终回到 chunk_api 去掉 vector 字段 → 字段名映射 → 返回前端渲染。

检索核心设计点

  1. 三段式混合检索:全文 query_string(字段加权)+ 稠密向量 knn + FusionExpr(weighted_sum, 0.05,0.95),融合权重在 ES 端用 1 - vector_similarity_weight 调全文 boost
  2. 向量不出引擎:ES 路径用二次 KNN-only 查询取干净 cosine,避免回传大向量
  3. 后端差异化:Infinity 自带 atan 归一无需本地 rerank;OceanBase 本地 rerank;ES 二次 KNN
  4. 最终打分:sim = tkweight*term + vtweight*vector + rank_feature(PageRank/标签)
  5. 召回兜底:首次空结果时降 min_match(0.3→0.1)、提 similarity(→0.17) 重试

关键文件清单

上传与存储链路

阶段文件关键函数/行号
前端 hookweb/src/hooks/use-document-request.tsuseUploadNextDocument:82
前端 serviceweb/src/services/knowledge-service.tsuploadDocument:321
后端路由api/apps/restful_apis/document_api.pyupload_document:369 / _upload_local_documents:571
存储工厂common/settings.pyStorageFactory:181
对象存储调用api/db/services/file_service.pyupload_document:512 / add_file_from_kb:430
MinIO 实现rag/utils/minio_conn.pyput:144
数据模型api/db/db_models.pyDocument:904 / File:934 / File2Document:949
文档服务api/db/services/document_service.pyinsert:440 / run:1044
地址反查api/db/services/file2document_service.pyget_storage_address:82
任务队列api/db/services/task_service.pyqueue_tasks:356
任务执行器rag/svr/task_executor.pybuild_chunks:265 / embedding:637 / do_handle_task:1447

检索链路

阶段文件关键函数/行号
前端表单web/src/pages/dataset/testing/testing-form.tsxformSchema:57
前端 hookweb/src/hooks/use-knowledge-request.tsuseTestRetrieval:62
前端 serviceweb/src/services/knowledge-service.tsretrievalTest:150
后端 APIapi/apps/restful_apis/chunk_api.pyretrieval_test:228
对话路径api/db/services/dialog_service.py:714
Query 改写rag/nlp/query.pyquestion:42 / hybrid_similarity:183
检索编排rag/nlp/search.pysearch:132 / retrieval:562 / get_vector:53
ES 存储rag/utils/es_conn.pysearch:141
Rerankrag/nlp/search.py / rag/llm/rerank_model.pyrerank_with_knn:443 / rerank_by_model:513