Spring Boot 3集成RocketMQ消息中间件:可靠消息发送与消费幂等实战

Spring Boot 3集成RocketMQ的消息架构选型

业务解耦和削峰填谷都要靠消息中间件。RocketMQ在金融、电商场景用得较多,延迟低、消息可达性有保障,且有事务消息。Spring Boot 3 + RocketMQ的组合,用官方starter即可快速接入,本文覆盖生产可用的发送、消费、事务消息与幂等方案。

项目依赖与RocketMQ生产端配置

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.3.0</version>
</dependency>

application.yml关键配置:

rocketmq:
  name-server: 192.168.1.20:9876
  producer:
    group: order-producer-group
    send-message-timeout: 3000

name-server是注册中心,broker把路由信息上报到这里,生产者和消费者都从它拿地址。

消息发送实战:同步发送与异步发送区别

同步发送在请求线程里等待broker确认,可靠但会占用线程时间;异步发送通过回调拿结果,吞吐更高:

@Service
public class OrderMessageService {

    @Resource
    private RocketMQTemplate rocketMQTemplate;

    // 同步发送,适合对延迟敏感的补偿/通知场景
    public SendResult sendOrderSync(Order order) {
        return rocketMQTemplate.syncSend(
            "order-topic",
            order,
            3000
        );
    }

    // 异步发送,适合量大的日志、推送类消息
    public void sendOrderAsync(Order order) {
        rocketMQTemplate.asyncSend("order-topic", order, new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) { }

            @Override
            public void onException(Throwable e) {
                // 落库待补偿,不要只打日志
                saveFailRecord(order, e.getMessage());
            }
        });
    }
}

发送失败必须处理:日志里补一条定时重试任务,比异步回调里”打印完不管”靠谱得多。

消息消费端配置:消费组与消息失败重试

@Component
@RocketMQMessageListener(
    topic = "order-topic",
    consumerGroup = "order-consumer-group"
)
public class OrderConsumer implements RocketMQListener<OrderMsg> {

    @Override
    public void onMessage(OrderMsg msg) {
        try {
            // 业务处理:创建单据、更新库存
            process(msg);
        } catch (Exception e) {
            // 抛异常触发重试,RocketMQ默认重试16次
            throw new RuntimeException("处理失败,等待重试", e);
        }
    }
}

消费者实现RocketMQListener接口,onMessage里抛异常即消费失败,broker按重试策略重新投递。注意消费端必须幂等:同一消息可能被重复消费,处理前先查业务状态或用去重表。

消息幂等与事务消息方案对比

RocketMQ的分布式事务消息用于本地事务和发消息必须一致的场景:先发半消息(half message),本地事务执行成功后commit,broker才把消息投递给消费者;失败则rollback。核心步骤:

@Service
public class OrderTxService {

    @Transactional
    public void createOrderWithTx(OrderDto dto) {
        // 1. 写订单表(本地事务)
        orderMapper.insert(dto);
        // 2. 半消息 + 本地事务状态上报
        rocketMQTemplate.sendMessageInTransaction(
            "order-tx-topic",
            MessageBuilder.withPayload(dto).build(),
            null
        );
    }

    @RocketMQTransactionListener
    class TxListener implements RocketMQLocalTransactionListener {
        @Override
        public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            // 本地事务已随createOrderWithTx提交,这里直接确认
            return LocalTransactionState.COMMIT_MESSAGE;
        }

        @Override
        public LocalTransactionState checkLocalTransaction(Message msg) {
            // broker回查:查数据库确认订单是否存在
            return orderExists(msg) ?
                COMMIT_MESSAGE : ROLLBACK_MESSAGE;
        }
    }
}

回查逻辑是兜底,必须查库确认事务结果,不能拍脑袋返回commit。事务消息只保证最终一致,业务上仍需幂等兜底。

RocketMQ常见故障排查与性能调优

几个高频问题的定位路径:

  • 消息发送超时:先看broker是否存活(mqadmin clusterList),再看网络与磁盘IO,同步发送超时建议调 send-message-timeout
  • 消费堆积:consumer并发度不足,调 consumerThreadMax 或扩容实例;确认消费链路没被外部调用拖慢
  • 重复消费:业务做去重(数据库唯一键、redis setnx),不要只依赖消息机制
  • 磁盘不足:broker默认48小时删除过期文件,空间紧张调 messageDelayLevel 与磁盘水位告警

上线前用压测把消费吞吐和堆积曲线摸清,再配置队列数(topic的queue数量决定并发上限,一般设16-32)。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot3-ji-cheng-rocketmq-xiao-xi-zhong-jian-jian-ke/

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

相关推荐