news 2026/10/2 6:11:42

第17章:RAGFlow API Server 与 Task Executor 双进程架构

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
第17章:RAGFlow API Server 与 Task Executor 双进程架构

1 项目背景

业务场景

「云帆科技」的知识库问答机器人上线两周后,运维小李发现了一个奇怪的现象:每天下午 2 点,当 HR 批量上传新一版的制度文件时,正在使用聊天功能的同事就会抱怨"回答好慢"“怎么转了半天没反应”。小李检查发现,文档上传后的解析任务和用户的聊天问答是在同一台机器的同一个 Python 进程中处理的——解析大 PDF 吃掉了大量 CPU,导致问答请求排队等待。

小李想起 RAGFlow 的架构介绍中提到过"API Server"和"Task Executor"是两个独立进程,但她一直用的是默认的一体化部署。这两个进程到底怎么分工?为什么要拆开?拆开后怎么调度?她决定深入理解这背后的架构设计。

痛点

不分拆进程的单体架构带来的问题:

  1. 资源争抢:CPU 密集的文档解析(OCR、Embedding)和 IO 密集的 HTTP 请求共用一个进程池,解析大文件时在线问答延迟飙升。
  2. 故障传播:解析一个损坏的 PDF 导致 Python 进程内存溢出,整个 RAGFlow 服务挂掉——连登录页面都打不开。
  3. 扩缩容困难:如果在线问答量大需要加实例,解析能力也跟着浪费式地扩容(反之亦然)。
  4. 灰度升级不便:想单独升级解析器版本,必须把整个服务停掉——在线问答也跟着断了。
单体架构 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 数量# ...

小白:“拆分后有什么收益?能具体量化吗?”

大师:“三大收益,每个都可以量化:”

  1. 故障隔离:Task Executor OOM 崩溃 → 不影响 API Server 响应 HTTP 请求 → 用户仍可查看已有知识库的问答(因为检索在 API Server 端完成)。
  2. 独立扩缩容:在线用户从 500 增到 2000 → 只扩容 API Server 副本数(2 → 4),Task Executor 保持 1 个即可。
  3. 独立升级:要升级文档解析器(如 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.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 异步★★★ 内存调用★★☆ 网络调用

适用场景

  1. 中等规模的 RAG 服务:在线问答 + 离线解析并存,需要隔离但不需要全微服务架构。
  2. 文档频繁更新的场景:每天有新文档批量导入,解析任务和在线问答需要物理隔离。
  3. SLA 要求高的服务:问答链路的可用性要求 99.9%,解析链路的偶尔中断可接受。
  4. 不同团队维护不同功能:后端团队负责 API,算法团队负责解析器——代码分离 + 进程分离。
  5. 灰度发布解析器:新版本解析器先在 Task Executor 上线,通过队列机制实现逐步切换。

不适用场景:

  1. 纯在线服务无离线任务:如果没有文档解析需求,Task Executor 多余。
  2. 解析和问答必须强一致:如果需要"上传完成瞬间即可检索",异步架构有延迟。

注意事项

  1. Redis 是单点故障:双进程通信完全依赖 Redis。如果 Redis 挂了,新上传的文档无法解析(API Server 正常但 Task Executor 饿死)。
  2. 数据库是共享状态:虽然进程分离,但它们共享同一个 MySQL——Task Executor 更新文档状态,API Server 读取状态。数据库锁竞争仍需关注。
  3. 日志被割裂:一次文档的全生命周期日志分别存在 API Server 和 Task Executor 中,排查问题时要看两个地方。
  4. Worker 数的内存陷阱:WS=2意味着 2 个 Worker 各加载一份 Embedding 模型(如 BGE-Large 1.2GB × 2 = 2.4GB 内存),乘以 2 也翻倍。
  5. 健康检查区分: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环境变量

思考题

  1. 当前双进程架构中,Task Executor 崩溃后,Redis 队列中的任务会保留。但如果在解析过程中崩溃(执行了一半),该文档的状态是 “parsing” 而非 “waiting_parse”——Task Executor 重启后不会重新处理这个任务。请设计一个"孤儿任务检测与恢复"机制。

  2. 如果 API Server 需要支持 5000 并发用户,单个实例不够用。需要使用 nginx 做负载均衡。但由于 RAGFlow 的多轮对话依赖服务端 Session,同一用户的两轮对话如果被分发到不同的 API Server 实例,第二轮的"追问"会丢失上下文。请设计负载均衡方案解决此问题。

(答案提示见第18章末尾或附录 D。)

延伸阅读与资源

10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/2 6:10:40

从零搭建AI工程能力:数据管道、模型训练到部署监控全流程实战

从零搭建AI工程能力这件事&#xff0c;我前前后后折腾过三回。第一回是跟着网上的教程跑通了几个Demo&#xff0c;觉得自己行了&#xff1b;第二回是接手一个真实项目&#xff0c;发现Demo和工程之间隔着一条河&#xff1b;第三回才算真正把整套东西理顺&#xff0c;从数据处理…

作者头像 李华
网站建设 2026/10/2 6:09:55

Linux网络编程进阶:数据边界、epoll事件驱动与线上排查实战

“Linux网络编程”这个系列能写到第四弹&#xff0c;说明前面的基础已经滚过了&#xff1a;socket 怎么创建、bind 和 listen 怎么配对、select 和 poll 怎么轮询、简单客户端服务端怎么跑通。按照我自己的习惯&#xff0c;到这一阶段就该换个视角了——不再问“这代码能不能跑…

作者头像 李华
网站建设 2026/10/2 6:09:53

OpenClaw本地安装实战:Node.js与Git环境准备及TaoToken接入配置

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/2 6:09:28

审稿----拒绝审稿的套话:用TaoToken统一Key跑通AI审稿工作流

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华