RagFlow 是一个基于深度文档理解的 RAG(检索增强生成)引擎。其核心数据流可分为两条主线:文件入库(上传→对象存储→解析→向量库)与检索召回(查询→混合检索→重排→返回)。
| 层 | 组件 | 职责 |
|---|---|---|
| 前端 | React + Vite + React Query | 表单采集、FormData 上传、检索参数提交 |
| API | Quart (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 / OceanBase | chunk 全文 + 稠密向量混合索引 ragflow_{tenant_id} |
api/apps/restful_apis/(非社区标准版的 document_app.py)。rag/utils/storage_factory.py 在本版本为空文件,真正的 StorageFactory 定义在 common/settings.py。
两条主线共享同一套存储与服务组件:实线为同步调用,虚线为异步消息流。
const restAPIv1 = `/api/v1`; documentUpload: (datasetId: string) => `${restAPIv1}/datasets/${datasetId}/documents`,
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);
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。
@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)
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)。
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]()
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)
@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 = kb.id(知识库 ID),key = location。两个装饰器实现单桶/多桶部署:单桶模式把 kb_id 拼进 key 形成 <prefix>/<kb_id>/<fnm>。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取消
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
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
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)
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)
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)
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)
kb_id;object key = location(=文件名,重名追加 _)thumbnail_{doc_id}.png;chunk 图片以 chunk id 为 keyragflow_{tenant_id},分区/路由键 = kb_idxxhash64(content_with_weight + doc_id);文档去重哈希 = xxhash128(blob)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]);
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 }; }
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`
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))
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.retrieval(dialog_service.py:714-729),参数来自 dialog 配置,结果经 kb_prompt 拼成 LLM 上下文。
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
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, ...)
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 再查一次
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))
atan 归一化(每路预归一后再融合),上层无需本地 rerank;OceanBase 仍回传向量走本地 rerank。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, ...)
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
_knn_scores) 取干净 cosine,再本地与 term 相似度按用户权重加权。外部 reranker 时向量分由模型 rerank_mdl.similarity 给出。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 - vector_similarity_weight 调全文 boostsim = tkweight*term + vtweight*vector + rank_feature(PageRank/标签)| 阶段 | 文件 | 关键函数/行号 |
|---|---|---|
| 前端 hook | web/src/hooks/use-document-request.ts | useUploadNextDocument:82 |
| 前端 service | web/src/services/knowledge-service.ts | uploadDocument:321 |
| 后端路由 | api/apps/restful_apis/document_api.py | upload_document:369 / _upload_local_documents:571 |
| 存储工厂 | common/settings.py | StorageFactory:181 |
| 对象存储调用 | api/db/services/file_service.py | upload_document:512 / add_file_from_kb:430 |
| MinIO 实现 | rag/utils/minio_conn.py | put:144 |
| 数据模型 | api/db/db_models.py | Document:904 / File:934 / File2Document:949 |
| 文档服务 | api/db/services/document_service.py | insert:440 / run:1044 |
| 地址反查 | api/db/services/file2document_service.py | get_storage_address:82 |
| 任务队列 | api/db/services/task_service.py | queue_tasks:356 |
| 任务执行器 | rag/svr/task_executor.py | build_chunks:265 / embedding:637 / do_handle_task:1447 |
| 阶段 | 文件 | 关键函数/行号 |
|---|---|---|
| 前端表单 | web/src/pages/dataset/testing/testing-form.tsx | formSchema:57 |
| 前端 hook | web/src/hooks/use-knowledge-request.ts | useTestRetrieval:62 |
| 前端 service | web/src/services/knowledge-service.ts | retrievalTest:150 |
| 后端 API | api/apps/restful_apis/chunk_api.py | retrieval_test:228 |
| 对话路径 | api/db/services/dialog_service.py | :714 |
| Query 改写 | rag/nlp/query.py | question:42 / hybrid_similarity:183 |
| 检索编排 | rag/nlp/search.py | search:132 / retrieval:562 / get_vector:53 |
| ES 存储 | rag/utils/es_conn.py | search:141 |
| Rerank | rag/nlp/search.py / rag/llm/rerank_model.py | rerank_with_knn:443 / rerank_by_model:513 |