news 2026/10/2 16:34:19

第18章:RAGFlow Redis 队列与文档解析任务调度

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
第18章:RAGFlow Redis 队列与文档解析任务调度

1 项目背景

业务场景

「云帆科技」的第 16 章综合实战交付后,系统平稳运行了一个月。但周一早晨,HR 部门一次性上传了 30 份新版制度的 PDF,触发了意想不到的问题:前 5 份文档在 2 分钟内就解析完成了,但从第 6 份开始,所有的文档都卡在"等待解析"状态,持续了整整 40 分钟。

运维小李排查后发现:Task Executor 的 Worker 默认只有 1 个,30 份文档按 FIFO(先进先出)顺序排队处理。更麻烦的是,有一份 200 页的扫描 PDF 占用了 Worker 长达 25 分钟,后面的 24 份文档只能干等。小李意识到:默认的 FIFO 队列不适合这种"大小文档混合"的场景——就像超市结账,一个人买了一整车,后面拿一瓶水的也得排半小时。

痛点

简单的 FIFO 任务队列在复杂场景下的局限:

  1. 大队列阻塞:一个超长任务(200 页扫描件)卡住整个队列,后续轻量任务全部饥饿。
  2. 无优先级:紧急文档和普通文档一视同仁——总监上传的年报和实习生上传的餐补规定同等排队。
  3. 失败处理粗暴:解析失败的文档只是标记"失败",没有自动重试机制——需要人工一个个手动重试。
  4. 无幂等保证: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)。要实现自动重试,有两个关键点:”

  1. 幂等性:同一个文档解析两次不能产生重复切片。RAGFlow 的处理方式是——重试前先清理该文档已生成的 Chunk。
  2. 退避策略:不要立即重试(可能是暂时的网络抖动),采用指数退避:第 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.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★★★ 原生★☆☆ 分区★★★ 原生
运维成本★★★ 零额外成本★★☆ 需维护★☆☆ 高★★☆ 需维护

适用场景

  1. 中小规模文档解析:日均 100-1000 份文档的解析调度,Redis Stream 足够。
  2. 需要任务可恢复性:Worker 进程可能因 OOM 或其他原因崩溃,PEL 保证任务不丢。
  3. 多 Worker 并行解析:通过消费组实现并行处理,动态扩缩 Worker 数。
  4. 任务优先级场景:紧急文档(领导要看的年报)先处理,普通文档后处理。
  5. 与现有 Redis 基础设施整合:不需要引入新的中间件,降低运维复杂度。

不适用场景:

  1. 海量文档实时流处理:10 万+ 文档/天,需要 Kafka 级别的高吞吐和分区并行。
  2. 严格顺序依赖的任务:文档 A 必须解析完才能解析文档 B——Stream 消费者组随机分配,无法保证顺序。

注意事项

  1. Stream 大小限制:Redis Stream 数据存在内存中,如果积压数万条消息且每条含大 payload,可能撑爆 Redis 内存。建议设置MAXLEN(~10000)。
  2. PEL 积压:如果 Worker 频繁崩溃且消息未 XACK,PEL 会持续增长,影响性能。需定期清理僵尸消息。
  3. Consumer Group 初始化:启动时XGROUP CREATE需要MKSTREAM参数(如果 Stream 不存在则自动创建)。
  4. Retry Count 无限增长:如果任务反复失败(如文件确实损坏),retry_count 会不断增长但永远不成功。需要最大重试次数上限 + 死信队列兜底。
  5. 多 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保障

思考题

  1. RAGFlow 当前的任务分配是 Worker 主动拉取(pull)模式。如果要在大量空闲时段节省 Worker 资源,你如何设计一个"按需伸缩"方案——队列积压超过阈值自动增加 Worker,积压清零后自动缩减 Worker?

  2. 某文档解析任务执行到一半时被 XCLAIM 认领到另一个 Worker 重新执行。如果原 Worker 此时也完成了任务并 XACK,就会出现同一文档被解析两次。请设计一个分布式锁方案(基于 Redis)来保证解析任务的互斥性。

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

延伸阅读与资源

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

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

Codex 100个真实案例 - 用AI做日志可视化分析平台(ELK替代方案)

/* 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 16:31:05

openrig 统一配置管理:Claude Code 与 Codex 多模型接入实战

1. openrig 到底是个什么东西第一次看到 openrig 这个名字&#xff0c;很多人会以为是某个硬件外设或者开源机械臂项目。实际上&#xff0c;结合它周围出现的关键词——Claude Code、Codex、YAML、Node.js——可以判断&#xff0c;这是一个围绕 AI 编程助手做统一接入与配置管理…

作者头像 李华
网站建设 2026/10/2 16:30:47

Cursor学习-Java环境配置:用TaoToken统一Key打通settings.json与JDK

/* 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 16:28:36

Process Lasso进程调度/资源管理工具

Process Lasso&#xff08;圈内俗称小绿&#xff09;Bitsum公司出品的Windows进程调度/资源管理工具&#xff0c;相当于增强版自动化任务管理器&#xff0c;专门管控进程CPU优先级、核心分配&#xff0c;解决CPU满载时系统卡死、前台程序卡顿的问题。✅ 核心功能ProBalance&…

作者头像 李华