分布式事务最终一致性方案:本地消息表与事务消息选型实战

微服务架构下,跨服务的数据操作无法依赖数据库本地事务保证ACID特性。分布式事务最终一致性方案通过消息中间件实现跨服务数据同步,是高并发设计中的核心组件。本文对比本地消息表和事务消息两种方案,给出选型建议和代码实现。

分布式事务问题场景分析

典型场景:电商系统下单后需要同时操作订单服务(创建订单)和库存服务(扣减库存)。两个服务使用独立数据库,本地事务无法跨库。如果订单创建成功但库存扣减失败,会出现超卖问题。

// 问题代码 - 非事务性操作
@Transactional
public void createOrder(OrderDTO dto) {
    // 1. 本地:创建订单
    orderMapper.insert(dto);

    // 2. 远程:扣减库存(可能失败)
    inventoryFeignClient.deduct(dto.getProductId(), dto.getQuantity());
    // 如果这里抛异常,订单已创建但库存未扣减
}

强一致性方案(如Seata AT模式)通过全局锁实现跨库事务,但性能损耗大,不适合高并发场景。最终一致性方案保证数据最终一致,在可接受短暂不一致的业务场景中是更优选择。

本地消息表方案实现

本地消息表的核心思路:将消息发送操作和业务操作放在同一个数据库事务中,保证业务操作和消息记录要么同时成功要么同时失败。后台任务轮询消息表,将未发送的消息投递到消息中间件。

-- 消息表DDL
CREATE TABLE local_message (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    business_key VARCHAR(64) NOT NULL COMMENT '业务唯一键',
    topic VARCHAR(64) NOT NULL COMMENT '消息主题',
    message_body TEXT NOT NULL COMMENT '消息内容JSON',
    status TINYINT DEFAULT 0 COMMENT '0-待发送 1-已发送 2-失败',
    retry_count INT DEFAULT 0 COMMENT '重试次数',
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
    update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    INDEX idx_status_create (status, create_time)
) ENGINE=InnoDB;
@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;
    @Autowired
    private MessageMapper messageMapper;
    @Autowired
    private RocketMQTemplate mqTemplate;

    @Transactional(rollbackFor = Exception.class)
    public void createOrder(OrderDTO dto) {
        // 1. 创建订单
        Order order = new Order();
        order.setOrderId(dto.getOrderId());
        order.setProductId(dto.getProductId());
        order.setQuantity(dto.getQuantity());
        order.setStatus("CREATED");
        orderMapper.insert(order);

        // 2. 同事务写入消息表
        LocalMessage msg = new LocalMessage();
        msg.setBusinessKey(dto.getOrderId());
        msg.setTopic("inventory-deduct");
        msg.setMessageBody(JSON.toJSONString(dto));
        msg.setStatus(0);
        messageMapper.insert(msg);
        // 事务提交后,订单和消息记录同时落库
    }
}

// 消息投递定时任务
@Component
public class MessageSendTask {

    @Scheduled(fixedDelay = 5000)
    public void sendPendingMessages() {
        List<LocalMessage> messages = messageMapper.selectPending(100);
        for (LocalMessage msg : messages) {
            try {
                mqTemplate.convertAndSend(msg.getTopic(), msg.getMessageBody());
                messageMapper.updateStatus(msg.getId(), 1);
            } catch (Exception e) {
                int retry = msg.getRetryCount() + 1;
                if (retry >= 5) {
                    messageMapper.updateStatus(msg.getId(), 2);
                    // 告警通知
                } else {
                    messageMapper.updateRetry(msg.getId(), retry);
                }
            }
        }
    }
}

本地消息表方案的优势:实现简单,不依赖消息中间件的事务特性,可适配任何MQ。劣势:定时轮询有延迟(通常5-10秒),消息表数据持续增长需要定期归档。

RocketMQ事务消息方案实现

事务消息方案依赖RocketMQ的两阶段发送机制。先发送半消息(对消费者不可见),执行本地事务后根据结果提交或回滚半消息。如果本地事务执行超时,RocketMQ回查本地事务状态。

@Service
public class OrderTransactionService {

    @Autowired
    private RocketMQTemplate mqTemplate;
    @Autowired
    private OrderMapper orderMapper;

    public void createOrder(OrderDTO dto) {
        Message<String> message = MessageBuilder
            .withPayload(JSON.toJSONString(dto))
            .setHeader("businessKey", dto.getOrderId())
            .build();

        mqTemplate.sendMessageInTransaction(
            "inventory-deduct",
            message,
            dto
        );
    }

    @RocketMQTransactionListener
    public class OrderTransactionListener
            implements RocketMQLocalTransactionListener {

        @Override
        public RocketMQLocalTransactionState executeLocalTransaction(
                Message msg, Object arg) {
            OrderDTO dto = (OrderDTO) arg;
            try {
                Order order = new Order();
                order.setOrderId(dto.getOrderId());
                order.setProductId(dto.getProductId());
                order.setQuantity(dto.getQuantity());
                order.setStatus("CREATED");
                orderMapper.insert(order);
                return RocketMQLocalTransactionState.COMMIT;
            } catch (Exception e) {
                return RocketMQLocalTransactionState.ROLLBACK;
            }
        }

        @Override
        public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
            String orderId = (String) msg.getHeaders().get("businessKey");
            Order order = orderMapper.selectByOrderId(orderId);
            if (order != null) {
                return RocketMQLocalTransactionState.COMMIT;
            }
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

事务消息方案的优势:消息投递实时性好(毫秒级),无需轮询,RocketMQ自带事务回查机制保证可靠性。劣势:强依赖RocketMQ,无法平滑切换到其他MQ。服务治理中需要额外维护事务回查接口。

两种方案选型对比与生产建议

选型决策矩阵:
维度          本地消息表      事务消息
消息延迟      5-10秒         毫秒级
MQ兼容性      任意MQ         仅RocketMQ
实现复杂度    中等           较高
数据库压力    有额外写       无额外表
事务回查      不支持         支持
运维成本      需归档清理     低

业务中台建设的实践建议:对消息实时性要求高的核心链路(支付、库存)使用事务消息方案;对实时性要求不高的旁路操作(通知、日志同步)使用本地消息表方案。API接口规范中应定义消息消费的幂等性要求,消费端通过business_key做去重。

// 消费端幂等处理
@RocketMQMessageListener(topic = "inventory-deduct")
public class InventoryConsumer implements RocketMQListener<OrderDTO> {

    @Autowired
    private InventoryMapper inventoryMapper;
    @Autowired
    private RedisTemplate redisTemplate;

    @Override
    @Transactional
    public void onMessage(OrderDTO dto) {
        String key = "deduct:" + dto.getOrderId();
        Boolean acquired = redisTemplate.opsForValue()
            .setIfAbsent(key, "1", 24, TimeUnit.HOURS);
        if (Boolean.FALSE.equals(acquired)) {
            return; // 已处理,跳过
        }

        int rows = inventoryMapper.deduct(
            dto.getProductId(), dto.getQuantity());
        if (rows == 0) {
            redisTemplate.delete(key);  // 回滚幂等标记
            throw new RuntimeException("库存不足");
        }
    }
}

幂等性是最终一致性方案的必要补充。消息中间件保证at-least-once投递语义,消费端必须处理重复消息。Redis SETNX配合业务唯一键是轻量级幂等方案,在高并发场景下性能表现稳定。

分布式事务最终一致性方案的选择没有银弹。理解业务对一致性和实时性的要求,结合团队能力和技术栈选型,才能做出合理的架构决策。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/fen-bu-shi-shi-wu-zui-zhong-yi-zhi-xing-fang-an-ben-di-xiao/

(0)
小编小编
上一篇 15小时前
下一篇 15小时前

相关推荐