1 项目背景
业务场景
「云帆科技」的知识库问答机器人上线两周后,运维小李发现了一个奇怪的现象:每天下午 2 点,当 HR 批量上传新一版的制度文件时,正在使用聊天功能的同事就会抱怨"回答好慢"“怎么转了半天没反应”。小李检查发现,文档上传后的解析任务和用户的聊天问答是在同一台机器的同一个 Python 进程中处理的——解析大 PDF 吃掉了大量 CPU,导致问答请求排队等待。
小李想起 RAGFlow 的架构介绍中提到过"API Server"和"Task Executor"是两个独立进程,但她一直用的是默认的一体化部署。这两个进程到底怎么分工?为什么要拆开?拆开后怎么调度?她决定深入理解这背后的架构设计。
痛点
不分拆进程的单体架构带来的问题:
- 资源争抢:CPU 密集的文档解析(OCR、Embedding)和 IO 密集的 HTTP 请求共用一个进程池,解析大文件时在线问答延迟飙升。
- 故障传播:解析一个损坏的 PDF 导致 Python 进程内存溢出,整个 RAGFlow 服务挂掉——连登录页面都打不开。
- 扩缩容困难:如果在线问答量大需要加实例,解析能力也跟着浪费式地扩容(反之亦然)。
- 灰度升级不便:想单独升级解析器版本,必须把整个服务停掉——在线问答也跟着断了。
单体架构 vs 双进程架构对比: 单体(不分家): [HTTP请求 + 文档解析] 混在同一个进程 解析10个PDF → CPU 100% → 用户聊天请求排队30秒 双进程(分家): [API Server] [Task Executor] 只处理HTTP请求 只处理文档解析 解析再忙也不影响聊天 聊天压力大也不影响解析2 项目设计
小胖:(指着监控图)“大师你看,每天下午 2 点 RAGFlow 的响应时间就飙到 20 秒!查了一下,都是 HR 在那个时间上传新文档导致的。这俩操作能不能互不影响?”
大师:“当然能,而且 RAGFlow 设计时就想到了这一点。你现在的部署是把 API Server 和 Task Executor 合在一起跑的,实际上它们应该是两个独立进程,就像餐厅的大堂和后厨——大堂只负责接待客人(API),后厨只负责做菜(解析),客人的点餐体验和后厨的忙碌程度互不影响。”
技术映射:API Server = 大堂服务员(处理点餐、传菜);Task Executor = 后厨团队(洗菜、切菜、炒菜)。两者各忙各的,通过传菜铃(Redis)沟通。
小胖:“那具体怎么分工的?什么请求走 API Server,什么请求走 Task Executor?”
大师:“分工非常清晰。API Server 负责所有在线、实时、同步的请求;Task Executor 负责所有离线、异步、耗时的任务。”
API Server 的职责: ├── HTTP 路由与请求分发 ├── 用户鉴权(JWT/API Token/Session) ├── 数据集 CRUD(创建、查看、修改、删除) ├── 文件上传接收(存储到 MinIO,创建 Document 记录) ├── 聊天问答(检索+LLM 生成,在线同步返回) ├── 文档状态查询(从数据库读,不直接查 Task Executor) └── 系统配置管理(模型供应商、用户、权限) Task Executor 的职责: ├── 从 Redis 队列消费解析任务 ├── PDF 解析(文本提取、OCR、版面分析) ├── 文档切片(Chunking) ├── Tokenizer 分词 ├── Embedding 向量化 ├── 写入文档引擎(ES/Infinity 索引构建) ├── 更新数据库中文档解析状态 └── GraphRAG 实体抽取与社区摘要(高级篇)小白:(在笔记本上画架构图)“那它们之间怎么通信?API Server 接收到上传文件后,怎么通知 Task Executor 来干活?”
大师:“唯一的通信渠道是 Redis。API Server 不直接调用 Task Executor,Task Executor 也不回掉 API Server——完全通过 Redis 队列解耦。”
API Server Redis Task Executor │ │ │ │ ──上传文件──▶ │ │ │ 1. 文件存MinIO │ │ │ 2. 数据库创建Document记录 │ │ │ 3. XADD ragflow_tasks │ │ │ {doc_id, action:parse} ──▶ │ │ │ │ ◀── XREADGROUP ── │ │ │ 拉取任务 │ │ │ │ │ │ 执行解析 │ │ │ Parser → Chunker │ │ │ → Embedding → Index │ │ │ │ │ ◀── 轮询GET /documents/{id} ─ │ ◀── 更新DB状态 ───────── │ │ 获取解析结果 │ status: success │技术映射:Redis 队列 = 厨房传菜单——服务员把订单夹在传送带上(XADD),厨师从传送带取单(XREADGROUP),做好了把菜放窗口(更新数据库状态),服务员自己去窗口取(轮询查询)。
小胖:“那源码层面是怎么把两者分开的?我看docker-compose.yml里好像有两个 service?”
大师:“对。RAGFlow 的启动脚本docker/launch_backend_service.sh根据环境变量决定启动哪个角色:”
# docker/launch_backend_service.sh 中的关键逻辑if["$LIGHTEN"=="1"];then# 轻量模式:只启动 API Server(不启 Task Executor)execpython3 api/ragflow_server.pyelse# 完整模式:同时启动两个进程python3 api/ragflow_server.py&# API Server 后台python3 rag/svr/task_executor.py&# Task Executor 后台waitfi# docker-compose.yml 中的拆分部署:services: ragflow-server: environment: -LIGHTEN=1# 只启动 API Server# ...ragflow-task-executor: environment: -LIGHTEN=0# 只启动 Task Executor-WS=2# Worker 数量# ...小白:“拆分后有什么收益?能具体量化吗?”
大师:“三大收益,每个都可以量化:”
- 故障隔离:Task Executor OOM 崩溃 → 不影响 API Server 响应 HTTP 请求 → 用户仍可查看已有知识库的问答(因为检索在 API Server 端完成)。
- 独立扩缩容:在线用户从 500 增到 2000 → 只扩容 API Server 副本数(2 → 4),Task Executor 保持 1 个即可。
- 独立升级:要升级文档解析器(如 PaddleOCR 版本)→ 只重启 Task Executor,API Server 无感知。
技术映射:双进程架构 = 消防隔断门——一个房间着火不会烧到另一个房间,人员(请求)可以从未着火的通道安全撤离。
3 项目实战
环境准备
目标:分别启动 API Server 和 Task Executor,观察日志中各自承担的职责。
前提:已按第16章完成 Docker Compose 部署。
分步实现
步骤1:观察一体化模式下的日志
目标:在未拆分的部署中观察日志,理解两个角色在同一进程的行为。
# 查看当前部署模式dockerps|grepragflow# 如果只有 ragflow-server 一个容器,说明是一体化模式# 查看一体化日志中两个角色的输出dockerlogs ragflow-server--tail100|grep-E"API|task|parse|request"典型日志示例(一体化):
[API] POST /api/v1/chats/xxx/messages - 200 (1.2s) ← API Server 职责 [API] GET /api/v1/datasets - 200 (0.1s) ← API Server 职责 [TASK] Received task: doc_id=abc, action=parse ← Task Executor 职责 [TASK] Parsing abc.pdf with DeepDoc... ← Task Executor 职责 [TASK] Embedding 47 chunks... ← Task Executor 职责 [API] POST /api/v1/chats/xxx/messages - 200 (3.8s) ← 解析时 API 响应变慢步骤2:拆分为独立容器部署
目标:修改 Docker Compose 配置,将两个角色部署为独立容器。
# docker-compose-split.yml - 拆分部署配置services:ragflow-api:image:infiniflow/ragflow:v0.26.0container_name:ragflow-apienvironment:-LIGHTEN=1# 仅 API Server-HTTP_PORT=9380-MYSQL_HOST=mysql-REDIS_HOST=redis-MINIO_HOST=minio-DOC_ENGINE=infinity-LOG_LEVEL=INFOports:-"80:9380"depends_on:mysql:condition:service_healthyredis:condition:service_healthyrestart:unless-stoppedmem_limit:2gragflow-task-executor:image:infiniflow/ragflow:v0.26.0container_name:ragflow-task-executorenvironment:-LIGHTEN=0# 仅 Task Executor-WS=2# 2 个 Worker 并行解析-MYSQL_HOST=mysql-REDIS_HOST=redis-MINIO_HOST=minio-DOC_ENGINE=infinity-LOG_LEVEL=INFOdepends_on:mysql:condition:service_healthyredis:condition:service_healthyrestart:unless-stoppedmem_limit:4g# Task Executor 需要更多内存# 启动拆分部署dockercompose-fdocker-compose-split.yml up-d# 确认两个容器dockerps--format"table {{.Names}}\t{{.Status}}\t{{.Image}}"# NAMES STATUS IMAGE# ragflow-api Up 2 minutes infiniflow/ragflow:v0.26.0# ragflow-task-executor Up 2 minutes infiniflow/ragflow:v0.26.0# ragflow-mysql Up 2 hours mysql:8.0# ragflow-redis Up 2 hours redis:7.2步骤3:验证故障隔离
目标:模拟 Task Executor 崩溃,验证 API Server 不受影响。
# 终端1:持续调用 APIwhiletrue;doSTATUS=$(curl-s-o/dev/null-w"%{http_code}"http://localhost/api/v1/version)echo"$(date+%H:%M:%S)API状态:$STATUS"sleep1done# 终端2:模拟 Task Executor 崩溃dockerstop ragflow-task-executorecho"Task Executor 已停止"# 观察终端1的输出 —— API Server 应持续返回 200,不受影响# 输出示例:# 14:30:01 API状态: 200 ← Task Executor 停掉后# 14:30:02 API状态: 200 ← API Server 仍然正常# 14:30:03 API状态: 200 ← 完全不受影响# 恢复 Task Executordockerstart ragflow-task-executor# 反向验证:API Server 崩溃不影响已入队的解析任务# 上传一个文档,观察解析开始curl-XPOST http://localhost/api/v1/datasets/<ds_id>/documents\-H"Authorization: Bearer$TOKEN"\-F"file=@large_doc.pdf"# 立即停掉 API Serverdockerstop ragflow-api# 检查 Task Executor —— 它已经拉取了 Redis 队列中的任务,应该继续执行dockerlogs ragflow-task-executor--tail5# [TASK] Parsing large_doc.pdf... (API Server 已停,但仍继续解析)# [TASK] Embedding chunks...# [TASK] Document large_doc.pdf parsed successfully步骤4:独立升级与扩容
目标:演示如何单独升级 Task Executor 或扩容 API Server。
# 场景1:独立升级 Task Executor(升级 DeepDoc 解析器版本)dockerpull infiniflow/ragflow:v0.27.0dockerstop ragflow-task-executordockerrmragflow-task-executor# 用新镜像启动(API Server 保持旧版本不变)dockercompose-fdocker-compose-split.yml up-dragflow-task-executor# 场景2:扩容 API Server 应对用户增长# 增加副本数dockercompose-fdocker-compose-split.yml up-d--scaleragflow-api=3# 前端加 nginx 负载均衡...# 场景3:根据队列积压动态调整 Task Executor Worker 数QUEUE_LEN=$(dockerexecragflow-redis redis-cli XLEN ragflow_tasks)if["$QUEUE_LEN"-gt50];then# 临时增加 Workerdockerexecragflow-task-executorsh-c"export WS=5 && python3 rag/svr/task_executor.py &"fi步骤5:监控指标分离
目标:为两个进程建立各自的监控指标。
# step5_monitoring.py - 分离监控importrequestsimporttimeimportjsondefcollect_api_metrics():"""收集 API Server 指标"""metrics={"active_requests":0,"request_rate_per_min":0,"avg_latency_ms":0,"error_rate":0,"5xx_count":0,"4xx_count":0,}# 从 API Server 日志或 /metrics 端点获取# RAGFlow 的 /api/v1/version 是一个轻量端点,可作为健康检查start=time.time()r=requests.get("http://localhost/api/v1/version",timeout=5)metrics["health_check_latency_ms"]=(time.time()-start)*1000metrics["healthy"]=r.status_code==200returnmetricsdefcollect_task_executor_metrics():"""收集 Task Executor 指标"""importsubprocess metrics={"queue_length":0,"active_workers":0,"parse_success_rate":0,"avg_parse_time_seconds":0,}# 从 Redis 获取队列长度result=subprocess.run(["docker","exec","ragflow-redis","redis-cli","XLEN","ragflow_tasks"],capture_output=True,text=True)metrics["queue_length"]=int(result.stdout.strip()or0)# 从 Task Executor 日志统计result=subprocess.run(["docker","logs","ragflow-task-executor","--since","5m"],capture_output=True,text=True)logs=result.stdout metrics["success_count"]=logs.count("parsed successfully")metrics["error_count"]=logs.count("ERROR")returnmetrics# 运行监控print("=== RAGFlow 双进程健康状况 ===")api_metrics=collect_api_metrics()task_metrics=collect_task_executor_metrics()print(f"API Server:{'健康'ifapi_metrics['healthy']else'异常'}")print(f" 延迟:{api_metrics['health_check_latency_ms']:.0f}ms")print(f"Task Executor:")print(f" 队列积压:{task_metrics['queue_length']}个任务")print(f" 近5分钟成功:{task_metrics['success_count']}")print(f" 近5分钟错误:{task_metrics['error_count']}")iftask_metrics["error_count"]>0:print(" ⚠ Task Executor 存在错误,请排查日志")iftask_metrics["queue_length"]>50:print(" ⚠ 队列积压严重,建议增加 Worker 数量")测试验证
# test_dual_process.py - 双进程架构验证测试deftest_api_available_when_task_executor_down():"""验证 API Server 在 Task Executor 宕机后仍可用"""# 停止 Task Executorimportsubprocess subprocess.run(["docker","stop","ragflow-task-executor"])time.sleep(5)# API Server 应该仍然正常r=requests.get("http://localhost/api/v1/version",timeout=5)assertr.status_code==200# 恢复subprocess.run(["docker","start","ragflow-task-executor"])deftest_task_resumes_after_executor_restart():"""验证 Task Executor 重启后继续处理积压任务"""# 上传一个测试文档ds=rag.create_dataset(name="TEST-任务恢复")doc=ds.upload_document("test_docs/sample.pdf")# 立即重启 Task Executor(模拟崩溃)subprocess.run(["docker","restart","ragflow-task-executor"])time.sleep(15)# 检查文档最终是否解析成功doc=ds.get_document(doc.id)assertdoc.status=="success",f"解析未恢复,状态:{doc.status}"# 清理rag.delete_dataset(ds.id)deftest_scale_api_independently():"""验证 API Server 可独立扩容"""# 查看当前 API Server 副本数result=subprocess.run(["docker","ps","--filter","name=ragflow-api","--format","{{.ID}}"],capture_output=True,text=True)api_count=len(result.stdout.strip().split("\n"))assertapi_count>=1,f"API Server 副本数:{api_count}"# 注:实际扩容命令为 docker compose scale完整代码清单
| 路径 | 说明 |
|---|---|
api/ragflow_server.py | API Server 入口,Quart 应用工厂 |
rag/svr/task_executor.py | Task Executor 入口,任务消费调度 |
docker/launch_backend_service.sh | 启动脚本,根据 LIGHTEN 决定启动角色 |
docker/docker-compose.yml | 容器编排(含拆分配置) |
api/apps/__init__.py | Blueprint 动态注册与鉴权 |
4 项目总结
优点 & 缺点
| 维度 | 双进程架构 | 单体架构 | 微服务架构 |
|---|---|---|---|
| 故障隔离 | ★★★ 进程级隔离 | ★☆☆ 完全耦合 | ★★★ 最强隔离 |
| 部署复杂度 | ★★★ 两个进程 | ★★★ 一个进程 | ★★☆ 多服务 |
| 独立扩缩容 | ★★★ 按角色扩缩 | ★☆☆ 整体扩缩 | ★★★ 按服务扩缩 |
| 资源利用率 | ★★☆ 各有空闲 | ★★★ 池化共享 | ★★☆ 各服务独立 |
| 运维复杂度 | ★★☆ 两份日志 | ★★★ 一份日志 | ★☆☆ N 份日志 |
| 通信开销 | ★★★ Redis 异步 | ★★★ 内存调用 | ★★☆ 网络调用 |
适用场景
- 中等规模的 RAG 服务:在线问答 + 离线解析并存,需要隔离但不需要全微服务架构。
- 文档频繁更新的场景:每天有新文档批量导入,解析任务和在线问答需要物理隔离。
- SLA 要求高的服务:问答链路的可用性要求 99.9%,解析链路的偶尔中断可接受。
- 不同团队维护不同功能:后端团队负责 API,算法团队负责解析器——代码分离 + 进程分离。
- 灰度发布解析器:新版本解析器先在 Task Executor 上线,通过队列机制实现逐步切换。
不适用场景:
- 纯在线服务无离线任务:如果没有文档解析需求,Task Executor 多余。
- 解析和问答必须强一致:如果需要"上传完成瞬间即可检索",异步架构有延迟。
注意事项
- Redis 是单点故障:双进程通信完全依赖 Redis。如果 Redis 挂了,新上传的文档无法解析(API Server 正常但 Task Executor 饿死)。
- 数据库是共享状态:虽然进程分离,但它们共享同一个 MySQL——Task Executor 更新文档状态,API Server 读取状态。数据库锁竞争仍需关注。
- 日志被割裂:一次文档的全生命周期日志分别存在 API Server 和 Task Executor 中,排查问题时要看两个地方。
- Worker 数的内存陷阱:
WS=2意味着 2 个 Worker 各加载一份 Embedding 模型(如 BGE-Large 1.2GB × 2 = 2.4GB 内存),乘以 2 也翻倍。 - 健康检查区分:API Server 健康不代表服务完全健康——还要检查 Task Executor 是否存活、队列是否积压。
常见踩坑经验
| 故障现象 | 根因 | 解决方法 |
|---|---|---|
| API Server 正常但上传后永远"等待解析" | Task Executor 已挂但没告警 | 对 Task Executor 也配健康检查 + 告警 |
| 队列积压几十个任务但 Worker 数=1 | 未设置 WS 环境变量,默认就是 1 | 设置WS=3或更高 |
| 扩容 API Server 后部分请求失败 | nginx 负载均衡未配置 sticky session(多轮对话需要) | 配置ip_hash或sticky cookie |
| 升级 Task Executor 镜像后解析全部失败 | 新版本 Parser 输出格式变了,索引写入不兼容 | 先在测试环境验证,不要直接升生产 |
| 两个进程的日志时间戳对不上 | 容器时区不一致 | 统一设置TZ=Asia/Shanghai环境变量 |
思考题
当前双进程架构中,Task Executor 崩溃后,Redis 队列中的任务会保留。但如果在解析过程中崩溃(执行了一半),该文档的状态是 “parsing” 而非 “waiting_parse”——Task Executor 重启后不会重新处理这个任务。请设计一个"孤儿任务检测与恢复"机制。
如果 API Server 需要支持 5000 并发用户,单个实例不够用。需要使用 nginx 做负载均衡。但由于 RAGFlow 的多轮对话依赖服务端 Session,同一用户的两轮对话如果被分发到不同的 API Server 实例,第二轮的"追问"会丢失上下文。请设计负载均衡方案解决此问题。
(答案提示见第18章末尾或附录 D。)
延伸阅读与资源
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析