第九章RabbitMQ 消息队列与索引任务异步化
很多 RAG 项目第一版都会犯一个错误:
用户上传文档后,后端立刻解析、切片、向量化、写入 pgvector。看起来链路很顺:上传成功,索引也顺手做完。
但只要文档稍微大一点,问题马上出现:
上传接口等待几十秒前端一直转圈PDF 解析失败后只看到接口超时Embedding 服务一慢,上传接口跟着变慢并发上传一多,后端线程被耗光
所以这篇文章真正要解决的不是“RabbitMQ 怎么用”,而是一个更现实的问题:
文档上传后,索引任务到底应该怎么异步化?结论先放在前面:
上传接口只负责保存文件、创建任务、发送消息。解析、切片、Embedding、向量入库,全部交给后台 task-service 异步执行。
RabbitMQ 在这里不是为了炫技。
它的价值是把“用户正在等待的上传接口”和“耗时、易失败、需要重试的索引链路”拆开。
这一步做好了,RAG 知识库才不会在文档上传阶段先把自己卡死。

01 别把上传和向量化写在一个接口里
很多刚开始做 RAG 系统的人,容易把“上传文档”和“索引文档”写在同一个接口里。
伪流程大概是这样:
上传文件-> 保存文件-> 解析文档-> 文本切片-> 调 Embedding-> 写入向量库-> 返回成功
这个写法在 Demo 里能跑,但在真实系统里风险很高。
上传接口不能长时间阻塞
上传接口应该是一个轻量操作。
它要做的事情是:
接收文件保存原始文件写入 document_info创建索引任务返回文档 ID 和状态
只要这几步成功,用户就应该看到“文档已上传,正在索引中”。
后面的解析、切片、向量化,不应该阻塞用户等待。
索引过程失败点很多
索引不是一个单点操作,而是一组可能失败的步骤。
比如:
MinIO 下载失败PDF 解析失败文本为空切片参数异常Embedding 服务超时向量库写入失败数据库连接异常
这些问题都不应该直接把上传接口拖垮。
更合理的做法是把索引过程抽象成一个后台任务,用任务状态记录每一步的结果。
任务需要重试和补偿
RAG 平台里的索引任务需要可追踪、可重试、可补偿。
用户上传文档后,系统至少要回答三个问题:
这个文档现在索引到哪一步了?失败原因是什么?失败后能不能重新执行?
如果没有任务表,只靠日志排查,后期运维会非常痛苦。
所以本章的关键不是“把 RabbitMQ 跑起来”,而是设计一条能落地的索引任务链路。
02 一个 RAG 索引任务到底要保存哪些状态
index_task 是索引任务的主表。
它不保存大文本,也不保存文件内容,只保存“这次索引任务”的元信息和状态。
一个简化版表结构可以这样设计:
CREATE TABLE index_task (id BIGSERIAL PRIMARY KEY,document_id BIGINT NOT NULL,kb_id BIGINT NOT NULL,user_id BIGINT NOT NULL,status VARCHAR(32) NOT NULL,retry_count INT NOT NULL DEFAULT 0,max_retry_count INT NOT NULL DEFAULT 3,error_message TEXT,start_time TIMESTAMP,end_time TIMESTAMP,create_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP);
核心字段并不复杂,但每个字段都很重要。
document_id 用来关联 document_info。
kb_id 表示这个文档属于哪个知识库。
user_id 用来做权限校验、任务归属和后续审计。
status 是任务流转的核心。
retry_count 和 max_retry_count 控制重试次数。
error_message 记录失败原因,方便后台页面展示和排查。
start_time、end_time 用来统计执行耗时,也可以配合超时扫描。

03 状态机设计:后台任务必须能看懂
任务表真正好不好用,关键看状态机是否清晰。
本项目中可以把索引任务拆成下面几种状态:
WAITING 等待执行RUNNING 正在执行SUCCESS 执行成功FAILED 执行失败RETRYING 等待重试TIMEOUT 执行超时
WAITING
knowledge-service 创建任务后,任务进入 WAITING。
这表示任务已经写入数据库,并且准备投递到 RabbitMQ。
RUNNING
task-service 消费到消息后,如果判断任务可以执行,就把状态改成 RUNNING。
这个状态表示后台已经开始处理。
SUCCESS
文档下载、解析、切片、写入 chunk、后续向量入库全部完成后,任务改成 SUCCESS。
同时 document_info 的状态也要更新为已索引。
FAILED
任务执行过程中出现不可恢复异常,就进入 FAILED。
这里一定要记录 error_message。
否则后台页面只能看到“失败”,却不知道为什么失败。
RETRYING
当管理员或系统触发重试时,任务可以从 FAILED 或 TIMEOUT 进入 RETRYING。
重试不是直接改成 RUNNING。
更清晰的做法是先标记重试中,再重新投递 MQ 消息,由消费者重新接管。
TIMEOUT
如果任务长时间处于 RUNNING,说明消费者可能宕机、卡死,或者某个外部服务一直没有返回。
此时可以由定时任务扫描并标记为 TIMEOUT。
哪些流转应该禁止
状态机要避免随意跳转。
例如:
SUCCESS -> RUNNING 不允许FAILED -> SUCCESS 不允许直接改TIMEOUT -> SUCCESS 不允许直接改WAITING -> SUCCESS 不允许跳过执行
状态越明确,系统越容易排查问题。
04 RabbitMQ 在这里解决的不是消息,而是解耦
RabbitMQ 在 RAG 索引链路里解决的不是“怎么发消息”,而是:
怎么让上传接口不用等待索引完成。knowledge-service 不直接调用 task-service 的索引接口,而是发送一条消息。
task-service 消费消息后,在后台慢慢完成下载、解析、切片、Embedding 和入库。
这样做有三个直接收益:
上传接口能快速返回索引失败可以独立记录和重试后台消费速度可以按资源情况扩展
RabbitMQ 的概念不需要背很多,先把下面几个和本项目相关的角色对齐即可。
Producer
生产者,负责发送消息。
在本项目里,knowledge-service 就是生产者。
Consumer
消费者,负责接收和处理消息。
在本项目里,task-service 就是消费者。
Exchange
交换机,负责接收生产者发送的消息,并根据规则转发到队列。
本项目可以使用 direct exchange。
Queue
队列,真正存放消息的地方。
消费者监听队列,从队列里取消息执行。
Routing Key
路由键,用来决定消息投递到哪个队列。
比如索引文档任务可以使用:
rag.index.documentAck
Ack 是消费者告诉 RabbitMQ:“这条消息我处理完了。”
如果处理成功后不 ack,消息可能一直停留在未确认状态。
如果异常时盲目 nack 并 requeue,消息可能反复回到队列,造成死循环。

05 本地启动 RabbitMQ:先把异步链路跑起来
开发阶段可以用 Docker 启动 RabbitMQ。
建议使用带管理后台的镜像:
docker run -d \--name rabbitmq \-p 5672:5672 \-p 15672:15672 \-e RABBITMQ_DEFAULT_USER=guest \-e RABBITMQ_DEFAULT_PASS=guest \rabbitmq:3-management
两个端口的作用不同:
5672 应用程序连接端口15672 RabbitMQ 管理后台端口
启动后可以访问:
http://localhost:15672默认用户名和密码都是:
guest / guestSpring Boot 里引入 RabbitMQ 依赖:
org.springframework.bootspring-boot-starter-amqp
application.yml 可以这样配置:
spring:rabbitmq:host: localhostport: 5672username: guestpassword: guestlistener:simple:acknowledge-mode: manual
这里建议使用手动 ack。
因为索引任务不是简单打印日志,而是会写数据库、写 MinIO、写向量库。只有业务真正处理完,才应该确认消息。
06 消息体别塞大对象,只传任务 ID
RabbitMQ 消息不要塞大对象。
尤其不要把文件内容、解析后的文本、切片结果放进消息体里。
消息只需要携带几个 ID:
{”taskId”: 1,”documentId”: 10,”kbId”: 12,”userId”: 3}
为什么只传 ID?
因为 ID 足够让消费者从数据库里查到所有上下文。
这样做有几个好处:
消息体小,投递快不会因为大文本导致 MQ 压力过大数据库是最终事实来源任务重试时可以重新读取最新状态消息结构更稳定
配套常量可以统一管理:
public class RabbitConstants {public static final String INDEX_EXCHANGE = ”rag.index.exchange”;public static final String INDEX_QUEUE = ”rag.index.queue”;public static final String INDEX_ROUTING_KEY = ”rag.index.document”;}
消息类可以保持简洁:
public class IndexTaskMessage {private Long taskId;private Long documentId;private Long kbId;private Long userId;}
消息越干净,后续越容易维护。
07 上传接口只做 3 件事:保存文件、建任务、发消息
knowledge-service 的职责是创建任务,并把任务交给 RabbitMQ。
推荐流程是:
保存 document_info保存原始文件到 MinIO创建 index_task,状态为 WAITING更新 document_info 状态为 INDEXING发送 MQ 消息返回上传结果

为什么先创建任务再发消息
因为 RabbitMQ 消息里会携带 taskId。
如果先发消息再创建任务,消费者可能立刻收到消息,但数据库里还查不到任务记录。
正确顺序应该是:
先落库再发消息
这样即使消费者很快执行,也能根据 taskId 查到任务。
消息发送失败怎么办
这里有一个真实系统一定会遇到的问题:
任务已经创建成功,但 MQ 消息发送失败了怎么办?
最简单的处理方式是让任务继续保持 WAITING,然后通过后台补偿任务扫描出来重新投递。
也可以在发送失败时把任务标记为 FAILED,并记录错误。
但从可恢复性角度看,WAITING + 补偿扫描 更适合索引任务。
因为任务本身还没有执行失败,只是消息暂时没有投出去。
document_info 状态也要同步
用户不会直接看 index_task 表。
用户更关心文档状态。
所以创建任务时,document_info 也要改成:
INDEXING这样前端列表可以显示“索引中”。
08 消费者别急着执行:先校验任务和幂等锁
task-service 是索引任务的真正执行者。
消费到消息后,不要立刻开干。
它应该先做三件事:
抢 Redis 幂等锁查询 index_task检查任务状态是否允许执行
只有这三步通过,才进入真正的索引流程。

一个简化后的消费流程如下:
@RabbitListener(queues = RabbitConstants.INDEX_QUEUE)public void handleIndexTask(IndexTaskMessage message, Channel channel, Message rawMessage) {Long taskId = message.getTaskId();try {// 1. Redis 锁防止重复消费// 2. 查询 index_task 并校验状态// 3. 标记任务 RUNNING// 4. 从 MinIO 下载文档// 5. 解析、归一化、切片// 6. 删除旧 chunk 和旧向量// 7. 写入新的 chunk 和向量// 8. 更新 document_info 和 index_taskchannel.basicAck(rawMessage.getMessageProperties().getDeliveryTag(), false);} catch (Exception e) {// 记录失败原因,更新任务状态channel.basicAck(rawMessage.getMessageProperties().getDeliveryTag(), false);}}
这里有一个容易忽略的点:
失败后也可以 ack。
原因是失败状态已经写入数据库了,后续重试应该由业务流程触发,而不是让 RabbitMQ 盲目重复投递。
为什么要先拿 Redis 锁
RabbitMQ 消息在某些场景下可能被重复投递。
例如消费者刚处理完业务,但还没 ack 就宕机。
RabbitMQ 会认为这条消息没有确认,于是重新投递。
所以消费者要先拿锁:
lock:index-task:{taskId}获取成功,说明当前实例可以执行。
获取失败,说明别的实例正在处理,当前消费者直接 ack 即可。
为什么还要检查任务状态
Redis 锁只能解决并发执行问题,不能替代数据库状态检查。
消费者拿到锁后,还要检查任务是否处于可执行状态:
WAITINGRETRYING
如果任务已经是 SUCCESS,说明它已经处理过了。
这时不能再次执行。
为什么要删除旧 chunk 和旧向量
重试或重新索引时,旧数据可能已经写入了一部分。
如果不先删除旧数据,就可能出现重复 chunk 或重复向量。
所以索引前要清理旧数据:
删除 document_chunk删除 document_vector重新写入 chunk重新写入 vector
这一步让“重建索引”变成一个可重复执行的动作。
09 失败、重试、超时:异步任务的生命线
异步任务不能只考虑成功路径。
真正上线以后,失败处理才是稳定性的关键。

失败时要做什么
任务失败时至少要更新两张表。
index_task:
status = FAILEDerror_message = 失败原因end_time = 当前时间retry_count = retry_count + 1
document_info:
status = FAILED这样前端才能看到文档索引失败,后台也能看到具体原因。
哪些错误适合重试
不是所有错误都值得重试。
适合重试的通常是临时性问题:
MinIO 网络抖动数据库瞬时连接失败Embedding 服务超时向量库短暂不可用
不适合自动重试的通常是数据本身的问题:
文件格式不支持PDF 损坏解析后文本为空切片参数配置错误
自动重试要谨慎,否则会把一个必然失败的任务重复执行很多次。
手动重试
管理后台可以提供“重试索引”按钮。
但不是所有状态都允许重试。
建议只允许:
FAILEDTIMEOUT
手动重试流程:
校验任务状态retry_count + 1status 改为 RETRYING清空 error_message重新发送 MQ 消息
这里不要直接调用索引方法。
统一通过 MQ 进入消费流程,链路会更一致。
超时扫描
如果任务长时间停留在 RUNNING,一般有两类原因:
消费者进程异常退出外部服务长时间无响应
可以用定时任务扫描:
SELECT *FROM index_taskWHERE status = 'RUNNING'AND start_time < now() - interval '30 minutes';
扫描到以后,把任务改成 TIMEOUT,并记录错误原因。
后续可以由管理员手动重试。
10 重复消费不可避免,幂等必须提前设计
消息队列系统里,必须默认消息可能重复。
所以消费者要做幂等设计。

第一层:Redis NX 锁
消费任务时,先用 Redis 加锁:
SET lock:index-task:{taskId} 1 NX EX 1800NX 表示不存在才设置。
EX 1800 表示锁 30 分钟后自动过期,避免进程异常退出后锁永远不释放。
第二层:任务状态检查
拿到锁以后,还要查数据库状态。
只有下面状态允许执行:
WAITINGRETRYING
如果已经是 SUCCESS、FAILED、TIMEOUT,就不应该继续执行。
第三层:数据库唯一约束
document_chunk 可以增加唯一约束:
ALTER TABLE document_chunkADD CONSTRAINT uk_document_chunk_indexUNIQUE (document_id, chunk_index);
这样即使代码层出现重复写入,数据库也能兜底拦住。
第四层:重建前删除旧数据
重试任务时,最容易出现旧数据残留。
所以重新索引前,要先删除旧 chunk 和旧向量。
顺序建议是:
删除旧向量删除旧 chunk重新写 chunk重新写向量
这几层加起来,才能让索引任务在重复消息、失败重试和服务重启后仍然保持稳定。
11 排查清单:任务卡住时先看这几处
下面这些问题,在调试 RabbitMQ 索引任务时很常见。
任务一直 WAITING
优先检查:
MQ 消息是否发送成功exchange / queue / routing key 是否一致task-service 是否启动消费者监听的队列名是否正确RabbitMQ 控制台里队列是否有消息堆积
如果任务在数据库里是 WAITING,队列里没有消息,重点查生产者发送逻辑。
如果队列里有消息但没人消费,重点查消费者服务。
任务一直 RUNNING
说明消费者已经开始执行,但没有正常结束。
重点检查:
task-service 日志MinIO 下载是否卡住PDF 解析是否卡住Embedding 调用是否超时数据库事务是否长时间未提交
这类任务要依赖超时扫描兜底。
任务 FAILED 但没有错误原因
说明异常捕获时没有把 error_message 写入数据库。
失败处理不能只打日志。
任务表里必须留下可展示、可排查的失败信息。
出现重复 chunk
重点检查三件事:
是否重复消费同一条消息重试前是否删除旧 chunkdocument_id + chunk_index 唯一约束是否存在
重复 chunk 会直接影响检索结果,必须尽早拦住。
消息反复重回队列
一般是消费者异常后执行了 nack requeue。
如果错误不可恢复,消息会不断被重新投递,造成循环失败。
索引任务更适合把失败写入数据库,然后 ack 消息。
后续由人工或补偿流程决定是否重试。
手动重试后旧消息仍在队列中
如果旧消息没有被正确 ack,管理员又触发了新消息,就可能出现同一个任务被执行两次。
解决方式仍然是幂等:
Redis 锁任务状态检查唯一约束重建前清理旧数据
消费速度过慢导致消息堆积
如果 RabbitMQ 队列里的消息不断增长,说明生产速度大于消费速度。
可以从几个方向优化:
增加 task-service 实例数调大消费者并发数优化 PDF 解析性能限制单个用户同时上传数量把 Embedding 调用做批量化
不要只盲目加消费者。
如果真正瓶颈在数据库或 Embedding 服务,加消费者反而可能把下游压垮。
12 动手验证:别只看上传成功,要看索引完成
这一章的验证目标不是只看接口返回成功。
你要把整条异步链路跑一遍。
第一步:启动 RabbitMQ
确认 Docker 容器正常运行。
访问管理后台:
http://localhost:15672检查 exchange、queue、binding 是否已经创建。
第二步:启动服务
至少需要启动:
knowledge-servicetask-serviceMinIOPostgreSQLRedisRabbitMQ
如果后续已经接入向量化,还要启动 Embedding 相关服务。
第三步:上传测试文档
通过前端或接口上传一个小文件。
建议先用 TXT 或 Markdown。
因为这类文件解析简单,方便先确认消息链路。
第四步:检查 index_task
上传后查询任务表。
你应该能看到一条新任务:
document_id = 当前文档 IDstatus = WAITING 或 RUNNING 或 SUCCESS
如果一直没有任务记录,说明任务创建逻辑有问题。
第五步:检查 RabbitMQ 队列
在 RabbitMQ 控制台查看:
rag.index.queue如果消息短暂出现后消失,说明消费者已经取走。
如果消息一直堆积,说明消费者没有正常消费。
第六步:检查 document_info
文档状态应该经历:
UPLOADED -> INDEXING -> INDEXED如果失败,则应该变成:
FAILED第七步:检查 document_chunk
索引成功后,查询 document_chunk。
应该能看到当前 document_id 对应的多条切片记录。
重点看:
chunk_index 是否连续content 是否为空token_count 或字符数是否合理
第八步:测试失败和重试
可以故意上传一个不支持的文件格式,或者临时关闭 MinIO,观察任务是否会进入 FAILED。
然后通过手动重试接口重新投递消息。
验证重点是:
状态是否从 FAILED / TIMEOUT 进入 RETRYING是否重新发送 MQ 消息是否最终成功或再次失败重复 chunk 是否被避免
本章小结
这一章的核心不是 RabbitMQ,而是 RAG 知识库里的一个工程边界:
上传接口不要承担索引任务。上传接口只应该做三件事:保存文件、创建 index_task、发送 MQ 消息,然后快速返回。
真正耗时、易失败、需要重试的解析、切片、Embedding 和向量入库,应该交给后台 task-service 异步执行。
index_task 负责把任务状态记录清楚,让系统知道文档现在是 WAITING、RUNNING、SUCCESS、FAILED、RETRYING 还是 TIMEOUT。
RabbitMQ 负责把 knowledge-service 和 task-service 解耦。消息体只携带 taskId、documentId、kbId、userId,不要把文件内容、解析文本和切片结果塞进 MQ。
任务消费时最重要的是三件事:
失败要能记录重试要能恢复重复消费不能写脏数据
Redis 锁、任务状态检查、数据库唯一约束、重建前清理旧 chunk / vector,是保证后台索引稳定运行的关键。
如果你的 RAG 项目现在还是“上传接口里直接解析、切片、向量化”,建议先把这一章收藏起来。
这个问题早期看不出来,等文档变大、并发变高、Embedding 服务变慢时,上传接口一定会先出问题。
下一篇,我们会继续进入 Embedding 与 pgvector,把已经切好的文本片段变成向量,并真正完成“知识库可检索”的核心能力。
作者有话说
如果这篇文章对你有帮助,欢迎点个关注。
这个专栏会持续更新KnowHub / RAG 平台实战内容,后面会继续把 Embedding、pgvector 向量检索、RAG 问答闭环和 Redis/Sentinel 等内容拆开讲清楚。
如果你想对照代码学习,可以结合下面两个仓库:
rag-demo-monolith:单体版源码
适合先理解 RAG 核心闭环,把业务链路跑通。
仓库地址:
https://gitee.com/MrLuoBin/rag-demo-monolith.git克隆命令:
git clone https://gitee.com/MrLuoBin/rag-demo-monolith.gitrag-platform:微服务版源码
适合继续学习 Gateway、Auth、Knowledge、Task 的企业级拆分方式。
仓库地址:
https://gitee.com/MrLuoBin/rag-platform.git克隆命令:
git clone https://gitee.com/MrLuoBin/rag-platform.git如果你正在做 AI 知识库或 RAG 平台,不要把“上传成功”误认为“文档可检索”。真正决定系统稳定性的,是后台索引任务能不能被追踪、失败后能不能恢复、重复执行时会不会写脏数据。
收藏和福利
觉得有用,可以转发给正在做 AI 知识库 / RAG 项目的队友。
建议先收藏,后面排查登录鉴权、文档上传、向量检索、RAG 问答和部署问题时,可以直接按章节回看。
也欢迎在评论区聊聊:你做 RAG 项目时,最头疼的是上传接口卡顿、索引任务失败、消息重复消费,还是 Embedding 调用太慢?
夜雨聆风