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 Stream(ragflow_tasks),Task 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 双持久 |
| 消息重试 | 需自行实现死信队列 | PEL(Pending 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 个 Worker:2 个处理 normal,1 个处理 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_count=0):# 第1步:查询当前文档状态doc=Document.get_by_id(doc_id)# 第2步:如果是重试,先清理旧的解析结果ifretry_count>0:# 删除已生成的 ChunkChunk.delete().where(Chunk.doc_id==doc_id).execute()# 删除文档引擎中的旧索引delete_index(doc.dataset_id,doc_id)# 第3步:执行解析try:chunks=parse_and_chunk(doc)embeddings=embed_chunks(chunks)save_to_engine(embeddings)# 第4步:更新状态并确认消息doc.status="success"doc.save()xack(task_id)# 确认消费exceptTemporaryErrorase:# 可重试的错误(网络超时、模型暂时不可用)ifretry_count<MAX_RETRIES:# 重新入队并设延迟retry_task(doc_id,retry_count+1,delay=exponential_backoff(retry_count))else:# 超限 → 死信队列move_to_dlq(doc_id,error=str(e))exceptPermanentErrorase:# 不可重试的错误(文件损坏、格式不支持)doc.status="failed"doc.error_msg=str(e)doc.save()xack(task_id)3 项目实战
环境准备
目标:上传 50 份文档,对比不同 Worker 数量下的解析吞吐和队列积压情况。
前提:RAGFlow 拆分部署(第17章),Task Executor 独立运行。
分步实现
步骤1:观察 Redis Stream 内部状态
目标:熟悉 Redis Stream 的监控命令。
# 连接 Redisdockerexec-itragflow-redis redis-cli# 查看所有 Stream>SCAN0TYPE 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 数量的吞吐量
目标:用 WS=1, WS=3, WS=5 三种配置跑同一批 30 个文档,对比解析耗时。
# 实验设计脚本#!/bin/bashDOC_COUNT=30forWSin135;doecho"=== 实验: WS=$WS==="# 设置 Worker 数并重启 Task Executordockercompose-fdocker-compose-split.yml up-d\--envWS=$WSragflow-task-executorsleep10# 等待启动# 记录开始时间START=$(date+%s)# 批量上传 30 个测试文档for((i=1;i<=DOC_COUNT;i++));docurl-XPOST"http://localhost/api/v1/datasets/$DS_ID/documents"\-H"Authorization: Bearer$TOKEN"\-F"file=@test_docs/sample_${i}.pdf"done# 等待全部解析完成(轮询)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, 处理中:$PENDING"if["$QUEUE_LEN"="0"];thenbreakfisleep5doneEND=$(date+%s)ELAPSED=$((END-START))THROUGHPUT=$(echo"scale=1;$DOC_COUNT/ ($ELAPSED/ 60)"|bc)echo" 完成! 总耗时:${ELAPSED}s (${THROUGHPUT}文档/分钟)"echo""done预期结果:
=== 实验: WS=1 === 完成! 总耗时: 840s (2.1 文档/分钟) === 实验: WS=3 === 完成! 总耗时: 310s (5.8 文档/分钟) === 实验: WS=5 === 完成! 总耗时: 210s (8.6 文档/分钟)坑点:Worker 数不是越多越好。5 个 Worker 各自加载 BGE-Large 模型 = 6GB 内存占用。如果机器内存只有 16GB,5 个 Worker 会触发 OOM Killer 随机杀死进程。
步骤3:模拟 Worker 崩溃与任务恢复
目标:验证 Stream 的 PEL 机制如何保证任务不丢失。
# 终端1:监控 PELwatch-n2'docker exec ragflow-redis redis-cli XPENDING ragflow_tasks task_executors'# 终端2:上传文档并在解析过程中杀死 Worker# 上传一个大文档,便于在解析过程中操作curl-XPOST"http://localhost/api/v1/datasets/$DS_ID/documents"\-H"Authorization: Bearer$TOKEN"\-F"file=@test_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_host="localhost",redis_port=6379):self.redis=redis.Redis(host=redis_host,port=redis_port)self.streams={"high":"ragflow_tasks:high","normal":"ragflow_tasks:normal","low":"ragflow_tasks:low",}self.group="task_executors"self.consumer_id=f"worker_{os.getpid()}"# 初始化 Stream 和消费组forstreaminself.streams.values():try:self.redis.xgroup_create(stream,self.group,id="0",mkstream=True)exceptredis.ResponseError:pass# 消费组已存在defenqueue(self,doc_id,priority="normal",metadata=None):"""入队任务"""stream=self.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,count=1,block_ms=5000):"""带优先级的消费:先取高优先级,再取普通,最后取低优先级"""# 按优先级顺序尝试消费forpriorityin["high","normal","low"]:stream=self.streams[priority]result=self.redis.xreadgroup(self.group,self.consumer_id,{stream:">"},count=count,block=0# 不阻塞,立即返回)ifresult:returnresultreturnNonedefacknowledge(self,stream,message_id):"""确认消息处理完成"""self.redis.xack(stream,self.group,message_id)defclaim_timeout_messages(self,min_idle_ms=60000):"""认领超时未确认的消息(故障恢复)"""forstreaminself.streams.values():pending=self.redis.xpending_range(stream,self.group,min="-",max="+",count=100)formsginpending:msg_id=msg["message_id"]idle_time=msg.get("time_since_delivered",0)ifidle_time>min_idle_ms:# 认领此超时消息claimed=self.redis.xclaim(stream,self.group,self.consumer_id,min_idle_ms,msg_id)ifclaimed:print(f"[RECOVERY] Claimed timeout message:{msg_id}")# 使用示例scheduler=PriorityTaskScheduler()# 入队不同优先级scheduler.enqueue("doc_urgent_report","high")scheduler.enqueue("doc_daily_update","normal")scheduler.enqueue("doc_archive_scan","low")# 消费(高优先级总是先被消费)whileTrue:tasks=scheduler.dequeue_with_priority()iftasks:forstream,messagesintasks:formsg_id,datainmessages:doc_id=data.get(b"doc_id",b"").decode()priority=data.get(b"priority",b"normal").decode()print(f"[PROCESS]{doc_id}(priority={priority})")# ... 执行解析 ...scheduler.acknowledge(stream,msg_id)else:time.sleep(1)步骤5:死信队列与失败告警
目标:建立失败任务的兜底机制。
# 创建死信队列监控脚本cat>dlq_monitor.sh<<'EOF' #!/bin/bash # RAGFlow 死信队列监控 DLQ_STREAM="ragflow_tasks:dlq" MAX_RETRIES=3 ALERT_WEBHOOK="https://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 EOFchmod+x dlq_monitor.sh测试验证
# test_task_scheduler.pyimportpytestimportredisimporttimeclassTestTaskQueue:deftest_fifo_order_preserved(self):"""验证 FIFO 顺序被保持"""r=redis.Redis(host="localhost",port=6379,db=0)stream="test_fifo_stream"r.delete(stream)# 入队 5 个任务foriinrange(5):r.xadd(stream,{"task_id":f"task_{i}","order":str(i)})# 消费并验证顺序messages=r.xread({stream:"0"},count=5)[0][1]orders=[m[1][b"order"].decode()for_,minmessages]assertorders==["0","1","2","3","4"]deftest_worker_recovery(self):"""验证 Worker 崩溃后消息可恢复"""r=redis.Redis(host="localhost",port=6379,db=0)stream="test_recovery_stream"group="test_group"r.delete(stream)try:r.xgroup_create(stream,group,id="0",mkstream=True)except:pass# 入队一个任务msg_id=r.xadd(stream,{"task":"important"})# Worker 1 消费但不确认(模拟崩溃)r.xreadgroup(group,"worker_dead",{stream:">"},count=1)# 确认消息在 PEL 中pending=r.xpending(stream,group)assertpending["pending"]==1# Worker 2 认领超时消息(min_idle=0 用于测试)claimed=r.xclaim(stream,group,"worker_recovery",0,[msg_id])assertlen(claimed)>0# 确认完成r.xack(stream,group,msg_id)pending=r.xpending(stream,group)assertpending["pending"]==0完整代码清单
| 路径 | 说明 |
|---|---|
rag/svr/task_executor.py | 任务消费主循环 + Worker 管理 |
rag/svr/task_queue.py | Redis Stream 队列操作封装 |
api/db/services/task_service.py | 任务状态与元数据管理 |
rag/flow/pipeline.py | 解析 Pipeline 触发入口 |
4 项目总结
优点 & 缺点
| 维度 | Redis Stream | RabbitMQ | Kafka | Celery + 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 频繁崩溃且消息未 XACK,PEL 会持续增长,影响性能。需定期清理僵尸消息。
- 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 实战修炼与源码剖析