
1 项目背景业务场景「云帆科技」的第 16 章综合实战交付后系统平稳运行了一个月。但周一早晨HR 部门一次性上传了 30 份新版制度的 PDF触发了意想不到的问题前 5 份文档在 2 分钟内就解析完成了但从第 6 份开始所有的文档都卡在等待解析状态持续了整整 40 分钟。运维小李排查后发现Task Executor 的 Worker 默认只有 1 个30 份文档按 FIFO先进先出顺序排队处理。更麻烦的是有一份 200 页的扫描 PDF 占用了 Worker 长达 25 分钟后面的 24 份文档只能干等。小李意识到默认的 FIFO 队列不适合这种大小文档混合的场景——就像超市结账一个人买了一整车后面拿一瓶水的也得排半小时。痛点简单的 FIFO 任务队列在复杂场景下的局限大队列阻塞一个超长任务200 页扫描件卡住整个队列后续轻量任务全部饥饿。无优先级紧急文档和普通文档一视同仁——总监上传的年报和实习生上传的餐补规定同等排队。失败处理粗暴解析失败的文档只是标记失败没有自动重试机制——需要人工一个个手动重试。无幂等保证Worker 处理到一半宕机重启后同一个任务可能被重新执行导致重复写入索引。FIFO 队列的长任务阻塞问题 Worker 1单线程处理顺序 [文档1: 2MB PDF, 50页] → 8分钟 [文档2: 200页扫描件] → 25分钟 ← 后面的全卡住 [文档3: 5KB Markdown] → 等25分钟 → 1秒完成 [文档4: 20KB DOCX] → 等25分钟 → 3秒完成 ... 后面24个文档平均等待时间25分钟2 项目设计小胖焦急地指着监控屏“大师你快看看Redis 队列里积压了 73 个任务Task Executor 只有一个 Worker 在吭哧吭哧干活后面的文档全在排队。我们能不能多开几个 Worker就像银行柜台——排队人多了多开几个窗口”大师“方向对了但没那么简单。多开 Worker多线程/多进程确实能并行处理多个文档但有几个问题需要考虑每个 Worker 会加载一份 Embedding 模型到内存2 个 Worker 2 倍内存不同 Worker 可能同时写入同一个文档引擎索引需要处理并发冲突。”技术映射Worker 银行柜台窗口开得越多同时服务的人越多但每个窗口都要配一个柜员内存而且多个柜员不能同时往同一个账户里存钱索引写入冲突。小胖“那 RAGFlow 具体是怎么做任务调度的按什么规则”大师“核心是 Redis Stream 消费组模型。RAGFlow 把每个解析任务包装成一条消息投递到 Redis Streamragflow_tasksTask Executor 的 Worker 作为消费组的消费者从 Stream 中拉取消息。”RAGFlow 任务调度架构 Redis Stream: ragflow_tasks ├── Consumer Group: task_executors │ ├── Worker-1 (consumer_id: worker_1_uuid) │ ├── Worker-2 (consumer_id: worker_2_uuid) │ └── Worker-3 (consumer_id: worker_3_uuid) │ 消息格式: { task_id: task_abc123, doc_id: doc_xyz789, dataset_id: ds_001, action: parse, priority: normal, // normal / high /紧急 retry_count: 0, created_at: 1718000000 }小白“Stream 和传统的 List 有什么区别为什么不用 LPUSH/BRPOP”大师“三个关键差异”特性Redis List (LPUSH/BRPOP)Redis Stream (XADD/XREADGROUP)消费确认弹出即删除无确认机制消费者显式 XACK未确认的消息可重新分配消费者组不支持原生支持多个消费者共享一个组消息回溯消费后不可回溯未确认的消息可重新读取消息持久化取决于 RDB/AOF 配置AOF RDB 双持久消息重试需自行实现死信队列PELPending Entries List天然支持“最重要的是——Stream 的消费者如果宕机它已经取走但未确认XACK的消息在 PEL 中其他消费者可以认领XCLAIM这些超时未确认的消息继续处理。这就是容错机制。”技术映射List 食堂打饭打完就没了Stream 快递签收系统——快递员取出包裹XREADGROUP但必须收件人签收XACK否则系统知道这件还没送达可以换个人送。小胖“那我如果想让重要文档优先处理怎么搞总不能改代码吧”大师“RAGFlow 在 Stream 消息中预留了priority字段。实现优先级队列的思路有多种”方案1多 Stream简单但粗粒度 ragflow_tasks:high → 高优先级 Stream ragflow_tasks:normal → 普通 Stream ragflow_tasks:low → 低优先级 Stream Worker 同时 XREADGROUP 三个 Stream优先处理 high 方案2加权轮询单 Stream对 Worker 逻辑简单 高优先级消息多分配 Consumer普通消息少分配 例如 3 个 Worker2 个处理 normal1 个处理 high 方案3消息优先级 XREADGROUP 排序依赖 Redis 7.2 Stream 消息可按 score 排序优先消费 score 高的小白“那解析失败的任务呢RAGFlow 有没有自动重试”大师“目前 RAGFlow 的默认行为是失败的任务不自动重试retry_count字段在消息中递增超过最大重试次数默认为 3后消息被移入死信队列DLQ。要实现自动重试有两个关键点”幂等性同一个文档解析两次不能产生重复切片。RAGFlow 的处理方式是——重试前先清理该文档已生成的 Chunk。退避策略不要立即重试可能是暂时的网络抖动采用指数退避第 1 次重试等 10 秒第 2 次等 30 秒第 3 次等 2 分钟。# 源码概念幂等性保证简化# 文件: rag/svr/task_executor.pydefhandle_parse_task(doc_id,retry_count0):# 第1步查询当前文档状态docDocument.get_by_id(doc_id)# 第2步如果是重试先清理旧的解析结果ifretry_count0:# 删除已生成的 ChunkChunk.delete().where(Chunk.doc_iddoc_id).execute()# 删除文档引擎中的旧索引delete_index(doc.dataset_id,doc_id)# 第3步执行解析try:chunksparse_and_chunk(doc)embeddingsembed_chunks(chunks)save_to_engine(embeddings)# 第4步更新状态并确认消息doc.statussuccessdoc.save()xack(task_id)# 确认消费exceptTemporaryErrorase:# 可重试的错误网络超时、模型暂时不可用ifretry_countMAX_RETRIES:# 重新入队并设延迟retry_task(doc_id,retry_count1,delayexponential_backoff(retry_count))else:# 超限 → 死信队列move_to_dlq(doc_id,errorstr(e))exceptPermanentErrorase:# 不可重试的错误文件损坏、格式不支持doc.statusfaileddoc.error_msgstr(e)doc.save()xack(task_id)3 项目实战环境准备目标上传 50 份文档对比不同 Worker 数量下的解析吞吐和队列积压情况。前提RAGFlow 拆分部署第17章Task Executor 独立运行。分步实现步骤1观察 Redis Stream 内部状态目标熟悉 Redis Stream 的监控命令。# 连接 Redisdockerexec-itragflow-redis redis-cli# 查看所有 StreamSCAN0TYPE stream# 返回包含 ragflow_tasks 的 key# 查看 Stream 基本信息XINFO STREAM ragflow_tasks# 返回: length消息总数, first-entry, last-entry, groups 等# 查看消费组XINFOGROUPSragflow_tasks# 返回: group: task_executors, consumers: 2, pending: 3# 查看待处理未确认的消息XPENDING ragflow_tasks task_executors# 返回: (待处理数量, 最早消息ID, 最晚消息ID, 各消费者分布)# 查看某个消费者未确认的消息列表XPENDING ragflow_tasks task_executors - 10预期输出示例XINFO STREAM ragflow_tasks length: 47 ← 当前积压 47 个任务 first-entry: 1718000000000-0 last-entry: 1718000045000-0 groups: 1 last-generated-id: 1718000045000-0 XPENDING ragflow_tasks task_executors 1) (integer) 3 ← 3 个任务正在处理中未确认 2) 1718000030000-0 ← 最早未确认 3) 1718000040000-0 ← 最晚未确认 4) 1) 1) worker_1 ← worker_1 有 2 个 2) 2 2) 1) worker_2 ← worker_2 有 1 个 2) 1步骤2对比不同 Worker 数量的吞吐量目标用 WS1, WS3, WS5 三种配置跑同一批 30 个文档对比解析耗时。# 实验设计脚本#!/bin/bashDOC_COUNT30forWSin135;doecho 实验: WS$WS# 设置 Worker 数并重启 Task Executordockercompose-fdocker-compose-split.yml up-d\--envWS$WSragflow-task-executorsleep10# 等待启动# 记录开始时间START$(date%s)# 批量上传 30 个测试文档for((i1;iDOC_COUNT;i));docurl-XPOSThttp://localhost/api/v1/datasets/$DS_ID/documents\-HAuthorization: Bearer$TOKEN\-Ffiletest_docs/sample_${i}.pdfdone# 等待全部解析完成轮询whiletrue;doQUEUE_LEN$(dockerexecragflow-redis redis-cli XLEN ragflow_tasks)PENDING$(dockerexecragflow-redis redis-cli XPENDING ragflow_tasks task_executors|head-1)echo 队列:$QUEUE_LEN, 处理中:$PENDINGif[$QUEUE_LEN0];thenbreakfisleep5doneEND$(date%s)ELAPSED$((END-START))THROUGHPUT$(echoscale1;$DOC_COUNT/ ($ELAPSED/ 60)|bc)echo 完成! 总耗时:${ELAPSED}s (${THROUGHPUT}文档/分钟)echodone预期结果 实验: WS1 完成! 总耗时: 840s (2.1 文档/分钟) 实验: WS3 完成! 总耗时: 310s (5.8 文档/分钟) 实验: WS5 完成! 总耗时: 210s (8.6 文档/分钟)坑点Worker 数不是越多越好。5 个 Worker 各自加载 BGE-Large 模型 6GB 内存占用。如果机器内存只有 16GB5 个 Worker 会触发 OOM Killer 随机杀死进程。步骤3模拟 Worker 崩溃与任务恢复目标验证 Stream 的 PEL 机制如何保证任务不丢失。# 终端1监控 PELwatch-n2docker exec ragflow-redis redis-cli XPENDING ragflow_tasks task_executors# 终端2上传文档并在解析过程中杀死 Worker# 上传一个大文档便于在解析过程中操作curl-XPOSThttp://localhost/api/v1/datasets/$DS_ID/documents\-HAuthorization: Bearer$TOKEN\-Ffiletest_docs/200page_scan.pdf# 等 10 秒解析开始sleep10# 强制杀死 Task Executordockerkillragflow-task-executor# 观察 PEL —— 应该有 1 条未确认的消息正在解析的那个文档# XPENDING 应该返回:# 1) (integer) 1# 重启 Task Executordockerstart ragflow-task-executorsleep15# 观察 PEL —— 重新启动后Worker 会从 PEL 中 Claim 超时的消息# 消息被重新处理最终确认PEL 变回 0步骤4实现优先级队列与长任务拆分目标扩展 RAGFlow 的任务调度逻辑区分紧急和普通任务。# extended_scheduler.py - 优先级调度演示概念代码importredisimportjsonimporttimeclassPriorityTaskScheduler:支持优先级的 RAGFlow 任务调度器def__init__(self,redis_hostlocalhost,redis_port6379):self.redisredis.Redis(hostredis_host,portredis_port)self.streams{high:ragflow_tasks:high,normal:ragflow_tasks:normal,low:ragflow_tasks:low,}self.grouptask_executorsself.consumer_idfworker_{os.getpid()}# 初始化 Stream 和消费组forstreaminself.streams.values():try:self.redis.xgroup_create(stream,self.group,id0,mkstreamTrue)exceptredis.ResponseError:pass# 消费组已存在defenqueue(self,doc_id,prioritynormal,metadataNone):入队任务streamself.streams.get(priority,self.streams[normal])message{doc_id:doc_id,priority:priority,retry_count:0,created_at:time.time(),metadata:json.dumps(metadataor{}),}self.redis.xadd(stream,message)print(f[ENQUEUE]{doc_id}-{priority}queue)defdequeue_with_priority(self,count1,block_ms5000):带优先级的消费先取高优先级再取普通最后取低优先级# 按优先级顺序尝试消费forpriorityin[high,normal,low]:streamself.streams[priority]resultself.redis.xreadgroup(self.group,self.consumer_id,{stream:},countcount,block0# 不阻塞立即返回)ifresult:returnresultreturnNonedefacknowledge(self,stream,message_id):确认消息处理完成self.redis.xack(stream,self.group,message_id)defclaim_timeout_messages(self,min_idle_ms60000):认领超时未确认的消息故障恢复forstreaminself.streams.values():pendingself.redis.xpending_range(stream,self.group,min-,max,count100)formsginpending:msg_idmsg[message_id]idle_timemsg.get(time_since_delivered,0)ifidle_timemin_idle_ms:# 认领此超时消息claimedself.redis.xclaim(stream,self.group,self.consumer_id,min_idle_ms,msg_id)ifclaimed:print(f[RECOVERY] Claimed timeout message:{msg_id})# 使用示例schedulerPriorityTaskScheduler()# 入队不同优先级scheduler.enqueue(doc_urgent_report,high)scheduler.enqueue(doc_daily_update,normal)scheduler.enqueue(doc_archive_scan,low)# 消费高优先级总是先被消费whileTrue:tasksscheduler.dequeue_with_priority()iftasks:forstream,messagesintasks:formsg_id,datainmessages:doc_iddata.get(bdoc_id,b).decode()prioritydata.get(bpriority,bnormal).decode()print(f[PROCESS]{doc_id}(priority{priority}))# ... 执行解析 ...scheduler.acknowledge(stream,msg_id)else:time.sleep(1)步骤5死信队列与失败告警目标建立失败任务的兜底机制。# 创建死信队列监控脚本catdlq_monitor.shEOF #!/bin/bash # RAGFlow 死信队列监控 DLQ_STREAMragflow_tasks:dlq MAX_RETRIES3 ALERT_WEBHOOKhttps://hooks.slack.com/xxx # 替换为实际告警 Webhook while true; do DLQ_COUNT$(docker exec ragflow-redis redis-cli XLEN $DLQ_STREAM 2/dev/null || echo 0) if [ $DLQ_COUNT -gt 0 ]; then echo [$(date)] ⚠ 死信队列中有 $DLQ_COUNT 个失败任务 # 获取死信详情 FAILED$(docker exec ragflow-redis redis-cli XRANGE $DLQ_STREAM - COUNT 5) # 发送告警 curl -s -X POST $ALERT_WEBHOOK \ -H Content-Type: application/json \ -d { \text\: \⚠ RAGFlow 死信队列告警\n死信数量: $DLQ_COUNT\n最近失败: $FAILED\ } fi sleep 300 # 每 5 分钟检查一次 done EOFchmodx dlq_monitor.sh测试验证# test_task_scheduler.pyimportpytestimportredisimporttimeclassTestTaskQueue:deftest_fifo_order_preserved(self):验证 FIFO 顺序被保持rredis.Redis(hostlocalhost,port6379,db0)streamtest_fifo_streamr.delete(stream)# 入队 5 个任务foriinrange(5):r.xadd(stream,{task_id:ftask_{i},order:str(i)})# 消费并验证顺序messagesr.xread({stream:0},count5)[0][1]orders[m[1][border].decode()for_,minmessages]assertorders[0,1,2,3,4]deftest_worker_recovery(self):验证 Worker 崩溃后消息可恢复rredis.Redis(hostlocalhost,port6379,db0)streamtest_recovery_streamgrouptest_groupr.delete(stream)try:r.xgroup_create(stream,group,id0,mkstreamTrue)except:pass# 入队一个任务msg_idr.xadd(stream,{task:important})# Worker 1 消费但不确认模拟崩溃r.xreadgroup(group,worker_dead,{stream:},count1)# 确认消息在 PEL 中pendingr.xpending(stream,group)assertpending[pending]1# Worker 2 认领超时消息min_idle0 用于测试claimedr.xclaim(stream,group,worker_recovery,0,[msg_id])assertlen(claimed)0# 确认完成r.xack(stream,group,msg_id)pendingr.xpending(stream,group)assertpending[pending]0完整代码清单路径说明rag/svr/task_executor.py任务消费主循环 Worker 管理rag/svr/task_queue.pyRedis Stream 队列操作封装api/db/services/task_service.py任务状态与元数据管理rag/flow/pipeline.py解析 Pipeline 触发入口4 项目总结优点 缺点维度Redis StreamRabbitMQKafkaCelery Redis部署复杂度★★★ 与 Redis 共用★★☆ 独立部署★☆☆ 重量级★★☆ pip install消息确认★★★ XACK 机制★★★ ACK/NACK★★☆ Offset commit★★☆ ACK消费者组★★★ 原生支持★★★ 原生★★★ 原生★★☆ 需配置吞吐量★★★ 10万/s★★☆ 中等★★★ 百万级★★☆ 中等优先级队列★★☆ 多 Stream★★★ 原生★☆☆ 分区★★★ 原生运维成本★★★ 零额外成本★★☆ 需维护★☆☆ 高★★☆ 需维护适用场景中小规模文档解析日均 100-1000 份文档的解析调度Redis Stream 足够。需要任务可恢复性Worker 进程可能因 OOM 或其他原因崩溃PEL 保证任务不丢。多 Worker 并行解析通过消费组实现并行处理动态扩缩 Worker 数。任务优先级场景紧急文档领导要看的年报先处理普通文档后处理。与现有 Redis 基础设施整合不需要引入新的中间件降低运维复杂度。不适用场景海量文档实时流处理10 万 文档/天需要 Kafka 级别的高吞吐和分区并行。严格顺序依赖的任务文档 A 必须解析完才能解析文档 B——Stream 消费者组随机分配无法保证顺序。注意事项Stream 大小限制Redis Stream 数据存在内存中如果积压数万条消息且每条含大 payload可能撑爆 Redis 内存。建议设置MAXLEN~10000。PEL 积压如果 Worker 频繁崩溃且消息未 XACKPEL 会持续增长影响性能。需定期清理僵尸消息。Consumer Group 初始化启动时XGROUP CREATE需要MKSTREAM参数如果 Stream 不存在则自动创建。Retry Count 无限增长如果任务反复失败如文件确实损坏retry_count 会不断增长但永远不成功。需要最大重试次数上限 死信队列兜底。多 Worker 的并发写入两个 Worker 同时往同一数据集索引写入 Chunk 时Infinity/ES 需要支持并发写入。常见踩坑经验故障现象根因解决方法Worker 数设为 5 但实际只有 1 个在工作消费者组中 consumer_id 相同导致被视为同一消费者每个 Worker 生成唯一 UUID 作为 consumer_id任务一直在队列中不被消费XREADGROUP 的符号未正确使用表示只读新消息确认代码中使用{stream: }而非{stream: 0}Redis 内存暴增Stream 中积累了 10 万 条未修剪的历史消息使用XADD ... MAXLEN ~ 10000自动修剪任务被重复处理两次Worker A 处理慢Worker B 通过 XCLAIM 认领了同一个任务在任务处理前检查 doc.status 是否已经是 “processing”消费组不存在错误Redis 重启后消费组信息丢失非持久化启动时用XGROUP CREATE的MKSTREAM保障思考题RAGFlow 当前的任务分配是 Worker 主动拉取pull模式。如果要在大量空闲时段节省 Worker 资源你如何设计一个按需伸缩方案——队列积压超过阈值自动增加 Worker积压清零后自动缩减 Worker某文档解析任务执行到一半时被 XCLAIM 认领到另一个 Worker 重新执行。如果原 Worker 此时也完成了任务并 XACK就会出现同一文档被解析两次。请设计一个分布式锁方案基于 Redis来保证解析任务的互斥性。答案提示见第19章末尾或附录 D。延伸阅读与资源10倍开发者的 Dify 魔法书从零构建全栈 AI 应用后端工程师转型AI第一课-Ollama 与私有化大模型实战大型语言模型(LLM) vLLM 高性能推理落地实战Agent开发之LlamaIndex 实战修炼与源码进阶大语言模型Transformers 实战修炼与源码剖析