Spring Boot整合RocketMQ实战:延迟消息与事务消息设计

RocketMQ是国内电商、金融场景使用最广的开源消息中间件之一。相比Kafka面向日志流,RocketMQ在消息可靠性、事务消息、延迟消息、顺序消息上都有成熟方案。本文用Spring Boot整合RocketMQ,从Producer、Consumer到延迟消息和事务消息的落地写法,并讨论消息幂等与堆积处理。

消息中间件选型与RocketMQ核心概念

RocketMQ的基础概念包括Topic(消息主题)、Tag(消息分类标签)、MessageQueue(队列,一个Topic下多个队列并行)、ConsumerGroup(消费组)。生产者发消息到Topic的某个Queue,消费者按组订阅,同一个Group内多个消费者实例分摊消息,不同Group独立消费。理解Queue模型是理解并行度与顺序性的关键:只有同一个Queue内的消息才保证顺序。

选型判断:需要事务消息、定时/延迟消息,RocketMQ(基于Java)提供;Kafka强在大吞吐流处理;RabbitMQ的复杂路由灵活性。业务系统里消息事务(订单创建+积分发放)和延迟消息(超时关闭订单)是两个高频需求。

Spring Boot整合RocketMQ:配置与生产者

# application.yml
rocketmq:
  name-server: 10.0.0.11:9876
  producer:
    group: order-producer-group
    send-message-timeout: 3000
  consumer:
    group: order-consumer-group
    topic: order-topic
// 生产者:注入RocketMQTemplate
@RestController
public class OrderController {
    @Resource
    private RocketMQTemplate rocketMQTemplate;

    @PostMapping("/order")
    public String createOrder(@RequestBody Order order) {
        order.setStatus("CREATED");
        // 同步发送,带Tag
        SendResult result = rocketMQTemplate.syncSend(
            "order-topic:CREATE", order, 5000);
        return "orderId=" + result.getMessageId();
    }
}

syncSend的第三个参数是超时时间,生产环境建议不传默认值,按业务容忍度显式设置。发送失败时业务侧要能降级处理:事务型操作不要依赖MQ作为唯一成功条件,失败消息用异步确认。

消费者:批量消费与幂等控制

@Service
@RocketMQMessageListener(
    topic = "order-topic",
    consumerGroup = "order-consumer-group",
    selectorExpression = "order",
    consumeMode = ConsumeMode.CONCURRENTLY
)
public class OrderConsumer implements RocketMQListener<Order> {
    @Override
    public void onMessage(Order order) {
        // 幂等:先查本地表,已处理过则跳过
        if (dedupService.isProcessed(order.getId())) return;
        try {
            orderService.updateStatus(order);
            dedupService.markProcessed(order.getId());
        } catch (Exception e) {
            // 抛出异常触发RocketMQ重试
            throw new RuntimeException("处理失败", e);
        }
    }
}

幂等是MQ消费的第一原则。RocketMQ重试机制默认每条消息重试16次,重试次数内抛异常会重新投递;消费失败超过次数进入死信队列(%DLQ%)。落库前先查去重表(消息ID或业务主键),保证即使重复投递也不产生重复业务操作。

延迟消息实现:超时订单场景

RocketMQ延迟消息不自由指定秒数,只支持18个预设等级:1s、5s、10s、30s、1m、2m、3m、4m、5m、6m、7m、8m、9m、10m、20m、30m、1h、2h。发消息时设置delayLevel,broker按等级延迟投递。

@Transactional
public void createOrderWithDelay(Order order) {
    // 创建订单
    orderMapper.insert(order);
    // 发送延迟消息:30分钟后检查支付状态
    Message<Order> msg = MessageBuilder.withPayload(order)
        .setHeader("delayLevel", 7)  // 7 = 10分钟
        .build();
    rocketMQTemplate.syncSend("order-topic:timeout", msg);
}

@RocketMQMessageListener(topic = "order-topic",
    consumerGroup = "order-timeout-group",
    selectorExpression = "timeout")
public class OrderTimeoutConsumer {
    public void onMessage(Order order) {
        if ("CREATED".equals(order.getStatus())) {
            // 未支付则关闭订单
            orderService.closeOrder(order.getId());
        }
    }
}

注意延迟消息的延迟等级由全局配置,不能动态指定秒数。要求更灵活的定时精度,可以用单独的调度任务扫描订单表,或用RocketMQ 5.x的新特性定时消息(时间戳格式)。

事务消息:本地事务与消息的原子性

典型场景:扣库存+发送MQ通知。两者要么都成功要么都回滚。事务消息流程:半消息发送到Broker,应用执行本地事务(扣库存SQL),提交结果通知Broker,Broker根据结果投递或回滚;异常时Broker回查本地事务状态,保底一致性。

// 发送事务消息
rocketMQTemplate.sendMessageInTransaction(
    "stock-topic:stock",
    MessageBuilder.withPayload(order).build(),
    order,  // arg 传给本地事务执行器
    new LocalTransactionListener() {
        @Override
        public LocalTransactionState executeLocalTransaction(
            Message msg, Object arg) {
            Order o = (Order) arg;
            try {
                stockService.deductStock(o);   // 本地事务
                return LocalTransactionState.COMMIT;
            } catch (Exception e) {
                return LocalTransactionState.ROLLBACK;
            }
        }
        @Override
        public LocalTransactionState checkLocalTransaction(
            Message msg) {
            // 回查:根据订单状态判断
            if (stockService.isStockDeducted(msg.getKeys())) {
                return LocalTransactionState.COMMIT;
            }
            return LocalTransactionState.ROLLBACK;
        }
    });

回查接口必须幂等,查询库确认事务结果。事务消息保证的是”本地事务+消息投递”的一致性,但消费侧的业务操作仍要幂等处理,因为消费者可能收到多次投递。

消息堆积与消费能力治理

堆积是MQ最常遇到的问题。排查方法:看Console的ConsumerDelay值,消息堆积总量、消费位点落后的Queue。应对手段:扩容消费者实例(注意实例数不能超过Queue数,否则部分实例空转)、拉长消费批量大小、拆分业务(消息里只放业务ID,处理逻辑异步化)。底层地按Queue消费,消息处理速度跟不上生产速度时,先用Redis合并同类业务,再批量写库,把高频小消息汇总成低频大批量消息。

RocketMQ的可靠性设计让它适合订单、支付这类强一致业务,把幂等、延迟、事务、堆积四个环节设计清楚,消息链路才能支撑高并发业务。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-zheng-he-rocketmq-shi-zhan-yan-chi-xiao-xi-yu/

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

相关推荐