数据一致性问题分析与解决方案
问题背景
最近学习es,做了个保单管理demo,修正了自己的一些错误认知 并且认识到将不可管理转换为可间接管理的模式
在保单管理系统中,PolicyService.createPolicy() 方法的执行流程如下:
PostgreSQL 保存(同步)
Redis 缓存(同步)
RabbitMQ 发送消息(同步)
Elasticsearch 同步(异步,通过 MQ)
这种多数据源协作模式带来了数据一致性的挑战。
核心问题
问题1:ES 同步失败导致的数据不一致
场景描述:
如果步骤 1-3 成功,但步骤 4(ES 同步)失败或延迟,会出现什么情况?
回答:
ES 中会丢失这条数据
导致复杂查询从 ES 中取不到数据,但简单查询却能从 DB 或缓存中获取
造成数据不一致
因为事务只在 1-3 步,MQ 中的消息消费处理和前 3 步没有绑定到一起
目前项目中没有处理这种情况
首次提出的解决方案:
分布式事务方案
使用 Seata AT 模式
通过全局事务来一起管理这 4 步的数据操作
但具体实现细节尚未考虑
验证和补偿机制
对 MQ 中发送成功的数据采用验证和补偿机制
确保消息成功消费
如果消费失败可以触发重试消费
或者进行异常通知
问题2:缓存成功但 DB 保存失败
场景描述:
如果 Redis 缓存成功,但 PostgreSQL 保存失败(事务回滚),缓存数据如何处理?
回答:
在当前代码中不会出现这种情况
因为 DB 异常会抛出异常,缓存没有机会执行
代码分析:
@Transactional
public Policy createPolicy(Policy policy) {
// 1. 保存到 PostgreSQL (主数据库)
Policy savedPolicy = policyRepository.save(policy);
// 2. 缓存到 Redis (提升查询性能)
redisTemplate.opsForValue().set(redisKey, savedPolicy, Duration.ofMinutes(30));
}
实际情况:
如果 PostgreSQL 保存失败:会抛出异常,Redis 缓存不会执行(用户理解正确)
如果 PostgreSQL 成功,但 Redis 失败:PostgreSQL 不会回滚(因为 Redis 不在 Spring 事务管理范围内)
问题3:MQ 消息发送失败
场景描述:
如果 MQ 消息发送失败,但 DB 和 Redis 已写入,如何保证 ES 最终同步?
回答:
在事务的管理下,DB 和缓存不会提交数据变更
这确保了数据一致性
MQ 有重试机制,在重试结束前,DB 和缓存在等待消息重试
如果全部失败,那么才会回滚
代码分析:
// 3. 发送消息到 RabbitMQ (异步处理)
try {
rabbitTemplate.convertAndSend(
RabbitMQConfig.POLICY_EXCHANGE,
RabbitMQConfig.POLICY_CREATED_KEY,
savedPolicy
);
System.out.println("[RabbitMQ] 已发送保单创建消息到 MQ");
} catch (Exception e) {
System.err.println("[RabbitMQ] 发送消息失败: " + e.getMessage());
}
实际情况(需要修正):
MQ 发送失败不会导致 DB 回滚
MQ 发送被 try-catch 包裹
即使发送失败,也不会抛出异常
事务不会回滚
convertAndSend 是异步的
方法立即返回,不等待 Broker 确认
不等待消息被消费
事务只管理 PostgreSQL
即使 MQ 发送失败,PostgreSQL 和 Redis 的操作已经完成
事务会正常提交
当前代码存在的风险
风险1:MQ 发送失败但数据已提交
// 当前代码:MQ 失败被捕获,不影响事务
try {
rabbitTemplate.convertAndSend(...);
} catch (Exception e) {
// 只打印日志,不抛异常
System.err.println(" [RabbitMQ] 发送消息失败");
}
// 事务继续提交,DB 和 Redis 数据已保存,但 MQ 消息丢失
影响:
DB 和 Redis 中有数据
但 ES 中没有数据(因为 MQ 消息丢失)
导致搜索功能无法找到该保单
风险2:Redis 不在事务范围内
@Transactional // 只管理 PostgreSQL
public Policy createPolicy(...) {
policyRepository.save(policy); // 在事务中
redisTemplate.set(...); // 不在事务中,失败不会回滚 DB
}
影响:
如果 Redis 操作失败,PostgreSQL 不会回滚
可能导致 DB 有数据但缓存没有
改进方案
方案1:本地消息表(推荐用于当前架构)
原理:
在同一个事务中,将需要发送的消息保存到本地数据库表,事务提交后,由定时任务扫描并发送消息。
实现步骤:
创建本地消息表
CREATE TABLE local_message (
id BIGSERIAL PRIMARY KEY,
policy_id VARCHAR(255) NOT NULL,
message_type VARCHAR(50) NOT NULL,
message_content TEXT,
status VARCHAR(20) DEFAULT 'PENDING',
retry_count INT DEFAULT 0,
created_at TIMESTAMP NOT NULL,
updated_at TIMESTAMP NOT NULL
);
修改业务代码
@Transactional
public Policy createPolicy(Policy policy) {
// 1. 保存到 DB
Policy savedPolicy = policyRepository.save(policy);
// 2. 保存到本地消息表(同一事务)
LocalMessage localMessage = new LocalMessage();
localMessage.setPolicyId(savedPolicy.getId());
localMessage.setMessageType("POLICY_CREATED");
localMessage.setStatus("PENDING");
localMessageRepository.save(localMessage); // 同一事务
// 3. 缓存到 Redis(失败不影响主流程)
try {
redisTemplate.opsForValue().set(...);
} catch (Exception e) {
// 记录日志,但不影响事务
}
return savedPolicy;
}
定时任务扫描并发送
@Scheduled(fixedDelay = 5000) // 每5秒执行一次
public void sendPendingMessages() {
List<LocalMessage> pendingMessages = localMessageRepository
.findByStatus("PENDING");
for (LocalMessage message : pendingMessages) {
try {
// 发送 MQ 消息
rabbitTemplate.convertAndSend(...);
// 更新状态为已发送
message.setStatus("SENT");
localMessageRepository.save(message);
} catch (Exception e) {
// 增加重试次数
message.setRetryCount(message.getRetryCount() + 1);
if (message.getRetryCount() >= 3) {
message.setStatus("FAILED");
}
localMessageRepository.save(message);
}
}
}
优点:
保证消息不丢失
实现简单,不需要额外中间件
可以记录消息发送状态
缺点:
需要额外的数据库表
有延迟(定时任务扫描间隔)
方案2:事务消息(需要 RabbitMQ 事务支持)
原理:
将 MQ 消息发送纳入事务管理,如果消息发送失败,回滚整个事务。
实现代码:
@Transactional
public Policy createPolicy(Policy policy) {
Policy savedPolicy = policyRepository.save(policy);
// 在事务中发送消息(需要配置事务管理器)
rabbitTemplate.execute(channel -> {
channel.txSelect(); // 开启事务
try {
rabbitTemplate.convertAndSend(...);
channel.txCommit(); // 提交
} catch (Exception e) {
channel.txRollback(); // 回滚
throw e; // 抛出异常,触发 Spring 事务回滚
}
return null;
});
return savedPolicy;
}
配置:
@Bean
public RabbitTransactionManager rabbitTransactionManager(
ConnectionFactory connectionFactory) {
return new RabbitTransactionManager(connectionFactory);
}
优点:
强一致性
消息发送失败会回滚整个事务
缺点:
性能较差(事务模式会降低吞吐量)
需要 RabbitMQ 支持事务模式
方案3:Seata AT 模式(分布式事务)
原理:
使用 Seata 的 AT 模式,通过全局事务协调器管理多个数据源的事务。
实现步骤:
引入依赖
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-alibaba-seata</artifactId>
</dependency>
配置 Seata
seata:
enabled: true
application-id: policy-service
tx-service-group: my_test_tx_group
config:
type: nacos
registry:
type: nacos
使用全局事务注解
@GlobalTransactional // Seata 全局事务
@Transactional
public Policy createPolicy(Policy policy) {
Policy savedPolicy = policyRepository.save(policy);
redisTemplate.opsForValue().set(...);
rabbitTemplate.convertAndSend(...);
return savedPolicy;
}
优点:
强一致性
支持多数据源事务管理
自动回滚机制
缺点:
需要部署 Seata Server
配置复杂
性能开销较大
方案对比
业务接受度分析
对于演示项目,当前的最终一致性是可以接受的
如果业务需要强一致性,建议使用 Seata 或本地消息表
状态机方案也是一个很好的思路,可以记录每个数据源的状态
总结
当前实现的特点
采用最终一致性模型
事务只管理 PostgreSQL
Redis 和 MQ 不在事务范围内
适合演示和学习场景
需要改进的地方
MQ 消息发送失败不会回滚事务
Redis 操作失败不会回滚 DB
缺少消息发送失败的补偿机制