
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 ExecutorTask 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[$LIGHTEN1];then# 轻量模式只启动 API Server不启 Task Executorexecpython3 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: -LIGHTEN1# 只启动 API Server# ...ragflow-task-executor: environment: -LIGHTEN0# 只启动 Task Executor-WS2# Worker 数量# ...小白“拆分后有什么收益能具体量化吗”大师“三大收益每个都可以量化”故障隔离Task Executor OOM 崩溃 → 不影响 API Server 响应 HTTP 请求 → 用户仍可查看已有知识库的问答因为检索在 API Server 端完成。独立扩缩容在线用户从 500 增到 2000 → 只扩容 API Server 副本数2 → 4Task Executor 保持 1 个即可。独立升级要升级文档解析器如 PaddleOCR 版本→ 只重启 Task ExecutorAPI Server 无感知。技术映射双进程架构 消防隔断门——一个房间着火不会烧到另一个房间人员请求可以从未着火的通道安全撤离。3 项目实战环境准备目标分别启动 API Server 和 Task Executor观察日志中各自承担的职责。前提已按第16章完成 Docker Compose 部署。分步实现步骤1观察一体化模式下的日志目标在未拆分的部署中观察日志理解两个角色在同一进程的行为。# 查看当前部署模式dockerps|grepragflow# 如果只有 ragflow-server 一个容器说明是一体化模式# 查看一体化日志中两个角色的输出dockerlogs ragflow-server--tail100|grep-EAPI|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_idabc, actionparse ← 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:-LIGHTEN1# 仅 API Server-HTTP_PORT9380-MYSQL_HOSTmysql-REDIS_HOSTredis-MINIO_HOSTminio-DOC_ENGINEinfinity-LOG_LEVELINFOports:-80:9380depends_on:mysql:condition:service_healthyredis:condition:service_healthyrestart:unless-stoppedmem_limit:2gragflow-task-executor:image:infiniflow/ragflow:v0.26.0container_name:ragflow-task-executorenvironment:-LIGHTEN0# 仅 Task Executor-WS2# 2 个 Worker 并行解析-MYSQL_HOSTmysql-REDIS_HOSTredis-MINIO_HOSTminio-DOC_ENGINEinfinity-LOG_LEVELINFOdepends_on:mysql:condition:service_healthyredis:condition:service_healthyrestart:unless-stoppedmem_limit:4g# Task Executor 需要更多内存# 启动拆分部署dockercompose-fdocker-compose-split.yml up-d# 确认两个容器dockerps--formattable {{.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状态:$STATUSsleep1done# 终端2模拟 Task Executor 崩溃dockerstop ragflow-task-executorechoTask 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\-HAuthorization: Bearer$TOKEN\-Ffilelarge_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-api3# 前端加 nginx 负载均衡...# 场景3根据队列积压动态调整 Task Executor Worker 数QUEUE_LEN$(dockerexecragflow-redis redis-cli XLEN ragflow_tasks)if[$QUEUE_LEN-gt50];then# 临时增加 Workerdockerexecragflow-task-executorsh-cexport WS5 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 是一个轻量端点可作为健康检查starttime.time()rrequests.get(http://localhost/api/v1/version,timeout5)metrics[health_check_latency_ms](time.time()-start)*1000metrics[healthy]r.status_code200returnmetricsdefcollect_task_executor_metrics():收集 Task Executor 指标importsubprocess metrics{queue_length:0,active_workers:0,parse_success_rate:0,avg_parse_time_seconds:0,}# 从 Redis 获取队列长度resultsubprocess.run([docker,exec,ragflow-redis,redis-cli,XLEN,ragflow_tasks],capture_outputTrue,textTrue)metrics[queue_length]int(result.stdout.strip()or0)# 从 Task Executor 日志统计resultsubprocess.run([docker,logs,ragflow-task-executor,--since,5m],capture_outputTrue,textTrue)logsresult.stdout metrics[success_count]logs.count(parsed successfully)metrics[error_count]logs.count(ERROR)returnmetrics# 运行监控print( RAGFlow 双进程健康状况 )api_metricscollect_api_metrics()task_metricscollect_task_executor_metrics()print(fAPI Server:{健康ifapi_metrics[healthy]else异常})print(f 延迟:{api_metrics[health_check_latency_ms]:.0f}ms)print(fTask 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 应该仍然正常rrequests.get(http://localhost/api/v1/version,timeout5)assertr.status_code200# 恢复subprocess.run([docker,start,ragflow-task-executor])deftest_task_resumes_after_executor_restart():验证 Task Executor 重启后继续处理积压任务# 上传一个测试文档dsrag.create_dataset(nameTEST-任务恢复)docds.upload_document(test_docs/sample.pdf)# 立即重启 Task Executor模拟崩溃subprocess.run([docker,restart,ragflow-task-executor])time.sleep(15)# 检查文档最终是否解析成功docds.get_document(doc.id)assertdoc.statussuccess,f解析未恢复状态:{doc.status}# 清理rag.delete_dataset(ds.id)deftest_scale_api_independently():验证 API Server 可独立扩容# 查看当前 API Server 副本数resultsubprocess.run([docker,ps,--filter,nameragflow-api,--format,{{.ID}}],capture_outputTrue,textTrue)api_countlen(result.stdout.strip().split(\n))assertapi_count1,fAPI Server 副本数:{api_count}# 注实际扩容命令为 docker compose scale完整代码清单路径说明api/ragflow_server.pyAPI Server 入口Quart 应用工厂rag/svr/task_executor.pyTask Executor 入口任务消费调度docker/launch_backend_service.sh启动脚本根据 LIGHTEN 决定启动角色docker/docker-compose.yml容器编排含拆分配置api/apps/__init__.pyBlueprint 动态注册与鉴权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 数的内存陷阱WS2意味着 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设置WS3或更高扩容 API Server 后部分请求失败nginx 负载均衡未配置 sticky session多轮对话需要配置ip_hash或sticky cookie升级 Task Executor 镜像后解析全部失败新版本 Parser 输出格式变了索引写入不兼容先在测试环境验证不要直接升生产两个进程的日志时间戳对不上容器时区不一致统一设置TZAsia/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 实战修炼与源码剖析