Administrator
发布于 2026-08-09 / 0 阅读
0
0

数据一致性问题分析与解决方案

数据一致性问题分析与解决方案

问题背景

最近学习es,做了个保单管理demo,修正了自己的一些错误认知 并且认识到将不可管理转换为可间接管理的模式
在保单管理系统中,PolicyService.createPolicy() 方法的执行流程如下:

  1. PostgreSQL 保存(同步)

  2. Redis 缓存(同步)

  3. RabbitMQ 发送消息(同步)

  4. Elasticsearch 同步(异步,通过 MQ)

这种多数据源协作模式带来了数据一致性的挑战。


核心问题

问题1:ES 同步失败导致的数据不一致

场景描述:
如果步骤 1-3 成功,但步骤 4(ES 同步)失败或延迟,会出现什么情况?

回答:

  • ES 中会丢失这条数据

  • 导致复杂查询从 ES 中取不到数据,但简单查询却能从 DB 或缓存中获取

  • 造成数据不一致

  • 因为事务只在 1-3 步,MQ 中的消息消费处理和前 3 步没有绑定到一起

  • 目前项目中没有处理这种情况

首次提出的解决方案:

  1. 分布式事务方案

    • 使用 Seata AT 模式

    • 通过全局事务来一起管理这 4 步的数据操作

    • 但具体实现细节尚未考虑

  2. 验证和补偿机制

    • 对 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());
}

实际情况(需要修正):

  1. MQ 发送失败不会导致 DB 回滚

    • MQ 发送被 try-catch 包裹

    • 即使发送失败,也不会抛出异常

    • 事务不会回滚

  2. convertAndSend 是异步的

    • 方法立即返回,不等待 Broker 确认

    • 不等待消息被消费

  3. 事务只管理 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:本地消息表(推荐用于当前架构)

原理:
在同一个事务中,将需要发送的消息保存到本地数据库表,事务提交后,由定时任务扫描并发送消息。

实现步骤:

  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
);
  1. 修改业务代码

@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;
}
  1. 定时任务扫描并发送

@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 模式,通过全局事务协调器管理多个数据源的事务。

实现步骤:

  1. 引入依赖

<dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-seata</artifactId>
</dependency>
  1. 配置 Seata

seata:
  enabled: true
  application-id: policy-service
  tx-service-group: my_test_tx_group
  config:
    type: nacos
  registry:
    type: nacos
  1. 使用全局事务注解

@GlobalTransactional  // Seata 全局事务
@Transactional
public Policy createPolicy(Policy policy) {
    Policy savedPolicy = policyRepository.save(policy);
    redisTemplate.opsForValue().set(...);
    rabbitTemplate.convertAndSend(...);
    return savedPolicy;
}

优点:

  • 强一致性

  • 支持多数据源事务管理

  • 自动回滚机制

缺点:

  • 需要部署 Seata Server

  • 配置复杂

  • 性能开销较大


方案对比

方案

一致性

复杂度

性能

适用场景

本地消息表

最终一致性

大多数业务场景

事务消息

强一致性

强一致性要求

Seata AT

强一致性

复杂分布式系统


业务接受度分析

  • 对于演示项目,当前的最终一致性是可以接受的

  • 如果业务需要强一致性,建议使用 Seata 或本地消息表

  • 状态机方案也是一个很好的思路,可以记录每个数据源的状态

总结

当前实现的特点

  • 采用最终一致性模型

  • 事务只管理 PostgreSQL

  • Redis 和 MQ 不在事务范围内

  • 适合演示和学习场景

需要改进的地方

  • MQ 消息发送失败不会回滚事务

  • Redis 操作失败不会回滚 DB

  • 缺少消息发送失败的补偿机制


评论