一个 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(); // 回滚 MySQL 写入
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, -- 'DOCUMENT'
aggregate_id VARCHAR(64) NOT NULL, -- 文档 ID
event_type VARCHAR(64) NOT NULL, -- 'DocumentUploaded'
payload JSON NOT NULL, -- 事件内容
status VARCHAR(16) NOT NULL DEFAULT 'PENDING', -- PENDING / PUBLISHED / FAILED
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) {
// 1. 写业务表
jdbc.update("INSERT INTO documents (id, title, content, owner) VALUES (?, ?, ?, ?)",
req.getId(), req.getTitle(), req.getContent(), req.getOwner());

// 2. 同事务写 outbox
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) // 每 1 秒扫一次
@Transactional
public void publishPending() {
// 1. 用 SELECT ... FOR UPDATE SKIP LOCKED 抢占
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 {
// 2. 投到 MQ
mq.asyncSend("rag.document.uploaded",
event.getPayload(),
new MessageQueueSelector() {
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object key) {
// 用 aggregate_id 做分区键,保证同一文档有序
int hash = Math.abs(key.hashCode() % mqs.size());
return mqs.get(hash);
}
},
event.getAggregateId() // 关键:分区键
);

// 3. 标记 published
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());
// 超过 N 次重试后改 FAILED,进死信
}
}
}
}

**核心机制是 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());

// 1. 幂等检查:向量表里是否已存在
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;
}

// 2. Spring AI 切分 + Embedding
List<Document> chunks = new TokenTextSplitter().split(
new DefaultContentFormatter().format(event.getContent())
);

List<float[]> embeddings = chunks.stream()
.map(chunk -> embeddingModel.embed(chunk.getContent()))
.toList();

// 3. 写入向量表(独立的、可重试的、不影响业务事务)
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 d
LEFT JOIN document_vectors v ON d.id = v.document_id
WHERE 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。一次性投入,长期受益。

关联阅读