1. 大文件并发场景下 RAG 的真实挑战
做过 RAG 项目的人大概都有个共识:Demo 跑通只要一个下午,但真把它丢到生产环境里,让几十上百号人同时上传几十兆甚至上百兆的文档,系统还能稳得住,那就是另一回事了。我前后经手过三四个 RAG 知识库项目,从最初用 LangChain 拼一个能问答的原型,到后来要支撑整个团队日常文档检索,中间踩的坑基本都集中在同一个地方——大文件 + 并发这两个条件叠加的时候,系统会以一种非常难看的方式崩掉。
这篇文章想聊的就是这件事:当你的 RAG 知识库需要同时处理大文件上传、解析、切分、向量化、入库这一整条链路,并且还要扛住多人并发的时候,到底该怎么设计。核心关键词就两个——RAG和大文件并发。我会把整条链路的拆解思路、关键参数怎么定、并发怎么控、瓶颈卡在哪、怎么排查,都摊开讲一遍。适合已经跑通过基础 RAG 流程、准备往生产环境推进的开发者,也适合正在被"上传一个大 PDF 就把服务卡死"折磨的同学。
先说清楚一个前提:RAG 本身不复杂,检索增强生成这套逻辑,无非是文档切块、向量化、存库、检索、拼上下文喂给大模型。真正难的是工程化。大文件并发之所以是个坎,是因为它同时压垮了三个环节——IO、CPU、内存,而且这三个环节的瓶颈会互相传导。你以为是向量库慢,其实是 PDF 解析把内存吃满了;你以为是模型推理慢,其实是切分逻辑在单线程里排队。不把这条链路拆清楚,调优就是瞎猜。
下面我按"整体设计思路 → 核心环节拆解 → 完整实操流程 → 问题排查"这个顺序来讲,每一部分都会给出我实际用过的方案和参数,能直接抄的我就写清楚,需要根据场景调整的我会说明判断依据。
2. 整体架构设计与并发思路拆解
2.1 为什么不能"上传即处理"
新手最容易犯的错,是把整个 RAG 流程写成同步的:用户上传文件 → 后端接收 → 解析 → 切分 → 向量化 → 入库 → 返回成功。小文件没问题,几百 KB 的 Markdown 秒级完成。但一个 80MB 的 PDF,光解析可能就要几十秒,向量化几千个 chunk 又要几分钟。这期间 HTTP 连接一直挂着,用户以为卡死了,服务器线程也被占着。来五个这样的请求,服务基本就废了。
所以第一个设计决策就是:上传和处理必须解耦。上传接口只负责把文件落到对象存储或本地磁盘,然后往消息队列里丢一个任务,立刻返回一个 task_id。真正的解析、切分、向量化交给后台的 worker 异步做。前端拿着 task_id 轮询进度,或者用 SSE/WebSocket 推状态。
这个决策背后的逻辑很直白:HTTP 请求的生命周期应该是短的,而文档处理的生命周期是长的,两者强行绑在一起,就是把长任务的成本转嫁给了连接层。解耦之后,上传接口的 QPS 可以很高,处理能力则通过 worker 数量独立伸缩。
2.2 并发模型的选择:进程、线程还是协程
Python 生态里做 RAG,绕不开 GIL。PDF 解析(比如 PyMuPDF)、文本切分、embedding 推理,这些都是 CPU 密集型操作,多线程在 GIL 下根本跑不满多核。我试过用 ThreadPoolExecutor 并发解析,实测下来 CPU 利用率卡在 120% 左右(假设 8 核),基本等于单核在跑。
所以并发模型我推荐多进程 + 队列的组合。用 Celery 或者更轻量的 RQ 做任务队列,worker 以进程方式启动,每个进程独立处理一个文档任务。进程数一般设为 CPU 核数,或者核数 + 1。这样每个进程能真正吃满一个核,8 核机器就能并行处理 8 个文档。
但这里有个坑:embedding 模型如果加载在每个 worker 进程里,8 个进程就是 8 份模型,内存直接爆炸。一个 bge-large 模型大概 1.3GB,8 份就是 10GB 起步。解决办法有两个:一是用独立的 embedding 服务(比如用 FastAPI 单独起一个服务,worker 通过 HTTP 调用),二是用共享内存或者模型量化。我一般选前者,虽然多了一次网络调用,但内存可控,而且 embedding 服务本身可以做批处理,吞吐反而更高。
2.3 大文件的分块处理策略
大文件不能一次性读进内存。一个 200MB 的 PDF,PyMuPDF 打开后如果一次性 extract 所有文本,内存峰值可能到 1GB 以上。并发几个这样的请求,内存就爆了。
我的做法是按页流式处理。PyMuPDF 支持逐页读取,处理完一页就释放一页的内存。具体来说,遍历每一页,提取文本,立刻做切分,切分后的 chunk 攒够一批(比如 64 个)就送去 embedding,embedding 完立刻入库,然后清空这批 chunk。这样内存占用是恒定的,跟文件大小无关,只跟批大小有关。
这个策略的关键在于:不要让整个文档的 chunk 同时存在于内存里。很多人习惯先把所有 chunk 切好存成一个 list,再统一 embedding,这在处理大文件时就是内存杀手。流式处理虽然代码复杂一点,但内存曲线是平的,非常稳。
2.4 向量库的写入并发
向量库这块,Milvus、Qdrant、Weaviate 我都用过。并发写入时最容易出问题的是批量大小和并发数的平衡。批量太小,网络往返次数多,吞吐上不去;批量太大,单次请求内存高,而且向量库服务端可能超时。
我的经验值是:单次批量 100 到 256 个向量,并发写入数控制在 4 到 8 之间。这个范围是实测出来的,再往上加,Qdrant 的写入延迟会明显上升,而且容易出现部分失败。另外,写入一定要做幂等,用文档 ID + chunk 序号做唯一键,重复写入时覆盖而不是追加,这样任务重试不会产生脏数据。
3. 核心环节拆解与关键参数实操
3.1 文档解析:不同格式的处理要点
RAG 知识库常见的文档格式有 PDF、Word、Markdown、HTML、纯文本。每种格式的解析坑都不一样。
PDF 是最麻烦的。扫描版 PDF 需要 OCR,这个成本很高,一般我会在解析前先判断:如果提取出的文本长度小于某个阈值(比如每页少于 50 个字符),就判定为扫描版,走 OCR 流程或者直接标记为"需人工处理"。文本版 PDF 用 PyMuPDF 就够了,速度快,内存可控。注意 PyMuPDF 的page.get_text()默认会保留一些排版信息,如果不需要,用get_text("text")更干净。
Word 文档用 python-docx,但要注意表格和图片。表格里的文本需要单独提取,否则会丢失。图片如果知识库需要支持图片检索(热词里有人问"rag知识库能存储图片嘛"),那就要走多模态 embedding,这是另一个话题,这里先不展开。
Markdown 和 HTML 相对简单,但要注意代码块。代码块里的内容如果被切分逻辑打散,检索时会出问题。我的做法是在切分前先识别代码块,把整个代码块作为一个不可分割的单元。
3.2 文本切分:chunk 大小和重叠的取舍
切分是 RAG 里最容易被忽视但影响最大的环节。chunk 太大,检索精度下降,因为一个 chunk 里混了太多主题;chunk 太小,上下文不完整,模型回答时缺信息。
我的默认参数是:chunk size 512 tokens,overlap 64 tokens。这个组合在大多数场景下表现均衡。但要注意,这是 token 数不是字符数。中文一个字符大约 1.5 到 2 个 token,所以 512 tokens 大概是 300 到 350 个中文字。如果你的文档是技术文档,句子长、术语多,可以适当放大到 768 tokens。
overlap 的作用是防止句子被切断导致语义丢失。64 tokens 的重叠大概能覆盖一到两句话,足够衔接上下文。overlap 太大会导致重复内容多,检索时返回一堆相似 chunk,浪费上下文窗口。
切分工具我用的是 LangChain 的 RecursiveCharacterTextSplitter,它支持按段落、句子、字符逐级降级切分,比单纯按字符数切要合理得多。配置的时候把 separators 设成["\n\n", "\n", "。", "!", "?", ".", " ", ""],中英文都能兼顾。
3.3 Embedding 批处理与并发控制
Embedding 是整条链路里最耗时的环节。一个 100 页的 PDF,切出来大概 500 到 800 个 chunk,用 bge-large 在 CPU 上跑,每个 chunk 大概 50 到 100ms,总共要 40 到 80 秒。如果用 GPU,能快 10 倍以上。
批处理是关键。单条 embedding 的效率极低,因为模型前向传播的固定开销占大头。把 batch size 设到 32 或 64,吞吐能提升 5 到 8 倍。但 batch 太大会吃内存,GPU 上尤其明显。我的经验是:GPU 显存 8GB 的话,bge-large 用 batch 64 没问题;CPU 的话,batch 32 比较稳。
并发控制这块,我建议用信号量限制同时进行的 embedding 请求数。如果是独立 embedding 服务,服务端自己会排队,客户端并发数设成 4 到 8 就行。如果是进程内调用模型,那并发数基本等于 worker 进程数,不需要额外控制,但要确保每个进程的 batch 不要太大,否则内存扛不住。
3.4 向量入库的批量与重试
入库这块,前面说了批量 100 到 256。这里补充一个细节:入库要和 embedding 流水线化。不要等所有 chunk 都 embedding 完再统一入库,而是 embedding 完一批就入库一批。这样即使中途失败,已经入库的部分不用重做,重试时从断点继续。
重试策略用指数退避,初始延迟 1 秒,最大延迟 30 秒,重试 3 次。如果 3 次还失败,把任务标记为失败,记录失败原因,人工介入。不要无限重试,否则一个坏任务会一直占着 worker。
4. 完整实操流程与关键代码
4.1 环境准备与依赖清单
先列一下我常用的技术栈,这套组合在多个项目里验证过,比较稳:
- 任务队列:Celery + Redis(Redis 同时做 broker 和 result backend)
- 文档解析:PyMuPDF(PDF)、python-docx(Word)、markdown(Markdown)
- 切分:LangChain 的 RecursiveCharacterTextSplitter
- Embedding:bge-large-zh-v1.5,用 FastAPI 单独起服务
- 向量库:Qdrant(轻量,单机部署方便)
- Web 框架:FastAPI
安装依赖:
pip install fastapi uvicorn celery redis pymupdf python-docx markdown langchain qdrant-client sentence-transformersQdrant 用 Docker 起一个单机版:
docker run -d -p 6333:6333 -v $(pwd)/qdrant_storage:/qdrant/storage qdrant/qdrant4.2 上传接口与任务分发
上传接口只做三件事:存文件、生成 task_id、丢任务。
from fastapi import FastAPI, UploadFile from celery import Celery import uuid, os app = FastAPI() celery_app = Celery('rag', broker='redis://localhost:6379/0', backend='redis://localhost:6379/1') UPLOAD_DIR = './uploads' os.makedirs(UPLOAD_DIR, exist_ok=True) @app.post('/upload') async def upload(file: UploadFile): task_id = str(uuid.uuid4()) file_path = os.path.join(UPLOAD_DIR, f'{task_id}_{file.filename}') # 流式写盘,避免大文件占内存 with open(file_path, 'wb') as f: while chunk := await file.read(1024 * 1024): f.write(chunk) celery_app.send_task('tasks.process_document', args=[task_id, file_path]) return {'task_id': task_id, 'status': 'queued'}注意这里用了await file.read(1024 * 1024)分块读,每次 1MB,避免把整个文件读进内存。这是处理大文件的基本功。
4.3 后台任务的流式处理实现
Celery 任务里做解析、切分、embedding、入库。核心是流式,不要攒。
from celery import shared_task import fitz from langchain.text_splitter import RecursiveCharacterTextSplitter import requests from qdrant_client import QdrantClient from qdrant_client.models import PointStruct @shared_task(bind=True, max_retries=3) def process_document(self, task_id, file_path): try: splitter = RecursiveCharacterTextSplitter( chunk_size=512, chunk_overlap=64, separators=["\n\n", "\n", "。", "!", "?", ".", " ", ""] ) qdrant = QdrantClient(host='localhost', port=6333) doc = fitz.open(file_path) batch_chunks, batch_ids = [], [] chunk_idx = 0 for page in doc: text = page.get_text("text") if len(text.strip()) < 50: continue for chunk in splitter.split_text(text): batch_chunks.append(chunk) batch_ids.append(f'{task_id}_{chunk_idx}') chunk_idx += 1 if len(batch_chunks) >= 64: flush_batch(qdrant, batch_chunks, batch_ids, task_id) batch_chunks, batch_ids = [], [] if batch_chunks: flush_batch(qdrant, batch_chunks, batch_ids, task_id) doc.close() return {'task_id': task_id, 'chunks': chunk_idx} except Exception as exc: raise self.retry(exc=exc, countdown=2 ** self.request.retries) def flush_batch(qdrant, chunks, ids, task_id): resp = requests.post('http://localhost:8001/embed', json={'texts': chunks}) vectors = resp.json()['vectors'] points = [ PointStruct(id=abs(hash(i)) % (10**12), vector=v, payload={'text': c, 'doc_id': task_id}) for i, c, v in zip(ids, chunks, vectors) ] qdrant.upsert(collection_name='docs', points=points)这段代码的关键点:每攒够 64 个 chunk 就 flush 一次,embedding 和入库都在 flush 里完成,然后清空列表。内存占用恒定在 64 个 chunk 的量级,跟文件大小无关。
4.4 Embedding 服务的批处理实现
独立的 embedding 服务,用 FastAPI 起,加载模型一次,常驻内存。
from fastapi import FastAPI from sentence_transformers import SentenceTransformer from pydantic import BaseModel app = FastAPI() model = SentenceTransformer('BAAI/bge-large-zh-v1.5') class EmbedRequest(BaseModel): texts: list[str] @app.post('/embed') def embed(req: EmbedRequest): vectors = model.encode(req.texts, batch_size=32, normalize_embeddings=True) return {'vectors': vectors.tolist()}normalize_embeddings=True很重要,归一化后的向量用余弦相似度检索时等价于内积,Qdrant 里可以用更快的距离计算方式。
4.5 并发 worker 的启动配置
Celery worker 用进程模式启动,并发数设为 CPU 核数:
celery -A tasks worker --loglevel=info --concurrency=8 --pool=prefork--pool=prefork是默认的多进程模式,每个 worker 进程独立处理任务。8 核机器设 8 个进程,每个进程处理一个文档,互不干扰。
但要注意,如果 embedding 是进程内调用(不是独立服务),那 8 个进程就是 8 份模型,内存扛不住。所以强烈建议用独立 embedding 服务,worker 只做解析、切分和 HTTP 调用,内存占用小很多。
5. 常见问题与排查技巧实录
5.1 内存暴涨的排查路径
内存暴涨是最常见的问题。排查顺序:先看是不是文件一次性读进了内存,再看是不是 chunk 列表攒太多,最后看是不是 embedding 模型重复加载。
用memory_profiler或者简单的psutil打点,在处理过程中每隔几秒记录一次 RSS。如果内存曲线是持续上升的,基本就是攒数据了;如果是阶梯式上升然后不降,可能是模型加载或者向量库客户端缓存。
我遇到过一次,内存一直涨,最后发现是 Qdrant 客户端的连接池没释放,每个任务创建一个新客户端,旧的不回收。改成全局单例客户端就好了。
5.2 任务卡死与超时处理
任务卡死通常是某个环节阻塞了。最常见的是 embedding 服务响应慢或者向量库写入超时。给每个 HTTP 调用加超时,embedding 设 30 秒,向量库写入设 10 秒。超时后抛异常,让 Celery 重试。
另外,Celery 要设task_time_limit和task_soft_time_limit,硬超时设 10 分钟,软超时设 8 分钟。软超时抛异常可以捕获,做清理;硬超时直接杀进程。没有超时限制的话,一个卡死的任务会永远占着 worker。
5.3 检索命中率低的调优
热词里有人提到 "rag hit rate",这确实是 RAG 的核心指标。命中率低通常是三个原因:chunk 切得不好、embedding 模型不适合领域、检索策略太单一。
chunk 问题前面说了,调整 size 和 overlap。embedding 模型如果领域特殊(比如医疗、法律),通用模型效果会差,需要微调或者换领域模型。检索策略上,单纯向量检索容易漏掉关键词匹配的情况,可以加一路 BM25 做混合检索,两路结果用 RRF(Reciprocal Rank Fusion)融合,命中率能提升 10 到 20 个百分点。
5.4 常见问题速查表
| 问题现象 | 可能原因 | 排查方法 | 解决方向 |
|---|---|---|---|
| 内存持续上涨 | chunk 列表攒太多 | 打点记录 RSS | 改流式处理,及时清空 |
| 任务卡死不结束 | 无超时限制 | 看 worker 日志 | 加 task_time_limit |
| 入库部分失败 | 批量太大或并发太高 | 看向量库日志 | 降批量到 128,并发降到 4 |
| 检索结果重复 | overlap 太大 | 检查 chunk 内容 | overlap 降到 32 到 64 |
| embedding 慢 | 单条调用或 batch 太小 | 看服务端 QPS | batch 提到 32 到 64 |
| 扫描版 PDF 无内容 | 未做 OCR | 检查提取文本长度 | 加 OCR 流程或标记人工 |
5.5 几个踩过的坑
第一个坑:Celery 的prefork模式下,如果在模块顶层加载模型,每个子进程都会加载一份,内存直接翻倍。解决办法是把模型加载放到函数内部,或者用worker_process_init信号在进程启动时加载一次。
第二个坑:Qdrant 的upsert如果 point id 用随机数,重复写入会产生重复数据。一定要用确定性的 id,比如文档 ID 加 chunk 序号的哈希。
第三个坑:大文件上传时,如果前端用FormData一次性提交,浏览器内存也会爆。前端也要做分片上传,后端合并。这个虽然不在 RAG 核心链路里,但实际项目中一定会遇到。
第四个坑:Redis 作为 broker 时,如果任务消息太大(比如把文件内容塞进消息里),Redis 内存会涨得很快。任务消息里只放文件路径和 task_id,不要放文件内容。
6. 性能压测与容量估算
6.1 单文档处理耗时拆解
以一个 50MB、200 页的文本版 PDF 为例,实测各环节耗时:
- 解析(PyMuPDF 逐页):约 8 秒
- 切分:约 2 秒
- Embedding(bge-large,GPU,batch 64):约 15 秒
- 入库(Qdrant,batch 128):约 5 秒
- 总计:约 30 秒
如果是 CPU 跑 embedding,这部分会变成 120 到 180 秒,总耗时 2 到 3 分钟。所以 GPU 对 RAG 的吞吐提升是决定性的。
6.2 并发容量估算
假设 8 核 CPU、32GB 内存、单张 8GB 显存的 GPU。worker 进程 8 个,每个进程处理一个文档,内存占用约 500MB(不含模型),8 个进程 4GB。embedding 服务占 2GB 显存和 3GB 内存。Qdrant 占 4GB 内存。总共约 11GB 内存,还有余量。
吞吐上,GPU embedding 是瓶颈。bge-large 在 8GB 显存上,batch 64 的吞吐大概是每秒 200 到 300 个 chunk。一个 200 页 PDF 约 600 个 chunk,需要 2 到 3 秒的 GPU 时间。8 个 worker 并发时,GPU 会排队,实际吞吐约每分钟 2 到 3 个文档。如果要更高吞吐,需要加 GPU 或者用更小的模型。
6.3 瓶颈定位方法
压测时用py-spy抓 worker 的调用栈,看时间花在哪。如果大量时间在requests.post,说明 embedding 服务是瓶颈;如果在qdrant.upsert,说明向量库是瓶颈;如果在page.get_text,说明解析是瓶颈。
定位到瓶颈后,针对性优化:embedding 瓶颈就加 GPU 或换小模型;向量库瓶颈就加索引或者分片;解析瓶颈就换更快的解析库或者预处理。
7. 一些延伸思考
这套方案跑下来,支撑团队内部几十人日常使用是没问题的。但如果要往更大规模走,还有几个方向可以扩展。
一是增量更新。现在每次上传都是全量处理,如果文档只是小改,全量重跑很浪费。可以做文档指纹(比如内容哈希),只处理变化的 chunk。这需要维护文档到 chunk 的映射关系,复杂度会上升,但能省很多算力。
二是多模态支持。热词里有人问 RAG 知识库能不能存图片。可以,但要用多模态 embedding 模型(比如 CLIP),把图片和文本映射到同一向量空间。这样检索时可以用文本搜图片,也可以用图片搜文本。实现上比纯文本复杂,主要是图片预处理和存储的成本。
三是GraphRAG 和本体 RAG。这两个是最近比较热的方向,核心思路是在向量检索之外,引入实体和关系图谱,解决"知识割裂"的问题。比如一个问题的答案分散在多个文档里,纯向量检索可能只召回其中一部分,而图谱能沿着关系把相关片段都找出来。这块我还在摸索,等有成熟经验了再单独写一篇。
四是检索质量评估。RAG 的 hit rate 不能靠感觉,要建评估集。人工标注一批问题和对应的正确 chunk,然后跑检索看召回率。有了评估集,调参才有方向,不然就是盲调。
这套东西说到底,核心就一句话:把长任务拆成短任务,把大内存拆成小内存,把串行拆成并行。大文件并发之所以难,是因为它同时挑战了这三个维度,任何一个没处理好,系统就会在压力下暴露问题。把链路拆清楚,每个环节的瓶颈定位准,参数调到位,剩下的就是工程耐心了。