一个 RAG 系统常见的”半成品”状态 你刚上线一个 RAG 知识库,流程是:
1 用户上传文档 → Spring AI 切分 + Embedding → 写入 MySQL 的向量表 → RAG 检索可用
跑通了。但 3 天后你发现一个诡异问题:
用户在管理后台明明看到”文档已上传”
但 RAG 问答时这条文档从来没被检索到
排查发现:MySQL 业务表写了,但向量表 INSERT 失败了——可能是 OOM、可能是事务冲突、可能是 Spring AI 的 batching 异常。结果就是业务数据”已上传”和向量数据”已索引”不一致 。
更糟糕的是:这种不一致默默存在 ,用户问不到,运营看不到,只有最细心的人翻日志才能发现。
这不是 Spring AI 的问题 这是分布式数据写入的经典问题 。任何”两个存储系统一起写”的场景都躲不开:MySQL + ES、MySQL + 缓存、MySQL + 搜索索引。RAG 只是最新的受害者。
三种”传统”解决方式都有明显缺陷:
方式 1:应用层 try-catch 回滚
1 2 3 4 5 6 7 8 9 transactionTemplate.execute(status -> { docRepository.save(doc); try { vectorRepository.save(embedding); } catch (Exception e) { status.setRollbackOnly(); throw e; } });
问题:向量表 INSERT 失败但 MySQL 已经回滚——下次重试时如果用户已经看到”上传失败”但实际没失败(事务边界没对齐),会重复入库。
方式 2:分布式事务(XA / Seata) 太重。一个 RAG 场景不值得引入 Seata 这套 AT 模式,运维成本陡增。
方式 3:先写向量后写业务表 更糟。向量表写成功了但业务表写失败,文档内容就凭空出现,没有归属信息。
正解:Outbox 模式 。
Outbox 模式的核心思想 不跨系统写两次。只在 MySQL 里写一次(业务数据 + outbox 记录一起写),由一个独立的 worker 扫 outbox 表,把消息投递到 MQ,下游消费者再去写向量表。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 Spring AI 服务 │ ├─→ MySQL 事务写入 │ ├─ business_documents(业务表) │ └─ outbox_events(事件表,状态=pending) │ ▼ 返回 200 OK(用户看到"上传成功") 异步:Outbox Worker 定时扫描 │ ├─ 取出 pending 事件 ├─ 投递到 Kafka/RocketMQ └─ 标记 outbox 事件状态 = published 异步:Vector Consumer 消费 │ ├─ 读 MQ 消息 ├─ Spring AI 生成 embedding └─ 写入 MySQL 向量表(独立的、最终一致)
关键不变量 :业务表和 outbox 表在同一个 MySQL 事务里写入 。要么都成功,要么都失败。不存在”业务表写成功但 outbox 写失败”的情况。
实现 Outbox 表 1 2 3 4 5 6 7 8 9 10 11 12 CREATE TABLE outbox_events ( id BIGINT PRIMARY KEY AUTO_INCREMENT, aggregate_type VARCHAR (64 ) NOT NULL , aggregate_id VARCHAR (64 ) NOT NULL , event_type VARCHAR (64 ) NOT NULL , payload JSON NOT NULL , status VARCHAR (16 ) NOT NULL DEFAULT 'PENDING' , created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP , published_at TIMESTAMP NULL , retry_count INT NOT NULL DEFAULT 0 , INDEX idx_status_created (status, created_at) );
注意 idx_status_created 这个索引——worker 扫描时就是 WHERE status='PENDING' ORDER BY created_at LIMIT 100,索引必须覆盖。
在 Spring AI 服务中写入 Outbox 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 @Service public class DocumentService { @Autowired private JdbcTemplate jdbc; @Transactional public void uploadDocument (DocumentUploadRequest req) { jdbc.update("INSERT INTO documents (id, title, content, owner) VALUES (?, ?, ?, ?)" , req.getId(), req.getTitle(), req.getContent(), req.getOwner()); String payload = objectMapper.writeValueAsString(Map.of( "documentId" , req.getId(), "title" , req.getTitle(), "content" , req.getContent() )); jdbc.update(""" INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload) VALUES (?, ?, ?, ?) """ , "DOCUMENT" , req.getId(), "DocumentUploaded" , payload); } }
业务事务只有 MySQL 一处 ,零跨系统复杂度。
Outbox Worker:扫表 + 投递 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 @Component public class OutboxPublisher { @Autowired private JdbcTemplate jdbc; @Autowired private RocketMQTemplate mq; @Scheduled(fixedDelay = 1000) @Transactional public void publishPending () { List<OutboxEvent> events = jdbc.query(""" SELECT * FROM outbox_events WHERE status = 'PENDING' ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED """ , new OutboxRowMapper ()); for (OutboxEvent event : events) { try { mq.asyncSend("rag.document.uploaded" , event.getPayload(), new MessageQueueSelector () { public MessageQueue select (List<MessageQueue> mqs, Message msg, Object key) { int hash = Math.abs(key.hashCode() % mqs.size()); return mqs.get(hash); } }, event.getAggregateId() ); jdbc.update("UPDATE outbox_events SET status='PUBLISHED', published_at=NOW() WHERE id=?" , event.getId()); } catch (Exception e) { jdbc.update("UPDATE outbox_events SET retry_count=retry_count+1 WHERE id=?" , event.getId()); } } } }
**核心机制是 SELECT ... FOR UPDATE SKIP LOCKED**:
FOR UPDATE 给这些行加行锁
SKIP LOCKED 让多个 worker 实例并发扫描时跳过别人正在处理的行
这样可以水平扩展 worker 数量,吞吐线性增长
消费者:MQ 消息 → 向量写入 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 @RocketMQMessageListener(topic = "rag.document.uploaded", consumerGroup = "vector-indexer") public class VectorIndexerConsumer implements RocketMQListener <MessageExt> { @Autowired private EmbeddingModel embeddingModel; @Autowired private JdbcTemplate jdbc; @Override public void onMessage (MessageExt msg) { DocumentEvent event = parse(msg.getBody()); Integer exists = jdbc.queryForObject( "SELECT COUNT(*) FROM document_vectors WHERE document_id = ?" , Integer.class, event.getDocumentId() ); if (exists > 0 ) { log.info("Vector already exists for doc {}, skip" , event.getDocumentId()); return ; } List<Document> chunks = new TokenTextSplitter ().split( new DefaultContentFormatter ().format(event.getContent()) ); List<float []> embeddings = chunks.stream() .map(chunk -> embeddingModel.embed(chunk.getContent())) .toList(); for (int i = 0 ; i < chunks.size(); i++) { jdbc.update(""" INSERT INTO document_vectors (document_id, chunk_index, content, embedding) VALUES (?, ?, ?, ?) """ , event.getDocumentId(), i, chunks.get(i).getContent(), serializeVector(embeddings.get(i))); } } }
失败处理与对账 Outbox 不是银弹。总有边界情况会失败 :
MQ 集群整个挂了,outbox 表里堆积了大量 PENDING
消费端代码 bug 持续失败,重试耗尽后进死信队列
向量表 INSERT 一直失败(主键冲突、字段超长、embedding 模型升级后维度变化)
3 道兜底防线 :
1. outbox 表里看堆积
1 2 3 SELECT status, COUNT (* ) FROM outbox_events WHERE created_at > NOW() - INTERVAL 1 HOUR GROUP BY status;
PENDING 数量超过阈值告警。
2. 死信队列人工处理 进入 DLQ 的消息不要自动重试,先看是什么问题。Embedding 模型升级导致维度不一致?修代码再回灌。
3. 定时全量对账(兜底) 每隔 N 小时跑一次:
1 2 3 4 SELECT d.id FROM documents dLEFT JOIN document_vectors v ON d.id = v.document_idWHERE v.id IS NULL AND d.created_at < NOW() - INTERVAL 1 HOUR ;
补 embedding、补入库。这是最后一道安全网。
为什么要这么做,而不是直接用 MQ 事务消息? RocketMQ 自带事务消息(Half Message + Commit/Rollback),看起来更”原生”。但有几个问题:
RocketMQ 事务消息的”回查”机制要求你的业务代码可重入 ——上传文档的逻辑必须幂等,实现复杂度反而上升
事务消息的”未决”状态有超时 (默认 1 分钟),如果业务逻辑跑超过这个时间,RocketMQ 会反查,反查逻辑写错就翻车
Outbox 模式对 MQ 不挑剔 ,换 Kafka、RabbitMQ 都不用改业务代码
Outbox 模式牺牲了一点点延迟(多了 worker 扫描的周期,默认 1 秒),换来了业务代码的简洁和 MQ 选型的自由度 。
意味着什么? 数据一致性问题在 AI 后端里比传统后端更严重 :
业务数据不一致 → 用户能看到、能投诉
RAG 数据不一致 → 用户问不到正确答案,但没任何错误提示 。这种”静默失败”对产品信任的伤害更大。
Outbox 模式不复杂,但能解决 90% 的”两个存储写不一致”问题。它是数据密集型 AI 后端工程师必会的模式 。
如果你的 RAG 系统还在用”双写 try-catch”或者”先写业务后写向量”,强烈建议改成 Outbox。一次性投入,长期受益。
关联阅读