RocketMQ 5.3.1 完整部署 + SpringCloud Alibaba 集成与消费幂等案例
一、安装部署1、安装包下载官方教程### 官方地址下载慢wget https://dist.apache.org/repos/dist/release/rocketmq/5.3.1/rocketmq-all-5.3.1-bin-release.zip### 可以手动去以下网址下载https://mirrors.aliyun.com/apache/rocketmq/5.3.1### 本文也提供了5.3.1的安装包2、linux部署1. 命令### 解压 unzip -x rocketmq-all-5.3.1-bin-release.zip 进入脚本目录 cd /home/rocketmq-all-5.3.1-bin-release/bin 启动mqnamesrv nohup sh mqnamesrv 等mqnamesrv 启动完成后再启动Broker nohup sh mqbroker -n localhost:9876 停止 先停Broker sh mqshutdown broker 再停NameServer sh mqshutdown namesrv NameServer 日志 tail -f /home/rocketmq-all-5.3.1-bin-release/logs/rocketmqlogs/namesrv.log Broker 日志 tail -f /home/rocketmq-all-5.3.1-bin-release/logs/rocketmqlogs/broker.log2、broker启动报错解决1.内存不够报错如下即内存不够默认需要8G内存Native memory allocation (mmap) failed to map 8589934592 bytes. There is insufficient memory for the Java Runtime Environment to continue. Command Line: -Xms8g -Xmx8g修改bin/runbroker.sh启动参数JAVA_OPT${JAVA_OPT} -server -Xms512m -Xmx512m -XX:MetaspaceSize128m -XX:MaxMetaspaceSize320m # 把15g改为1g JAVA_OPT${JAVA_OPT} -XX:MaxDirectMemorySize1g2.透明大页 THP 开启 alwaysBroker 直接退日志开头打印/sys/kernel/mm/transparent_hugepage/enabled: [always] madvise never原因:RocketMQ Broker 强制校验透明大页always 模式会造成内存碎片、刷盘卡顿启动拦截。# 临时生效服务器重启失效 echo madvise | tee /sys/kernel/mm/transparent_hugepage/enabled echo madvise | tee /sys/kernel/mm/transparent_hugepage/defrag # 校验 cat /sys/kernel/mm/transparent_hugepage/enabled # 备选跳过校验仅本地测试 nohup sh mqbroker -n 127.0.0.1:9876 -Drocketmq.broker.skipTransparentHugePageChecktrue 3. windows部署windows部署依旧跟linux用同一个zip包只是启动的命令不同而已启动前需要在环境变量中配置好JAVA_HOME。### 启动namesrv mqnamesrv.cmd 启动broker mqbroker.cmd -n 127.0.0.1:9876 autoCreateTopicEnabletrue二、集成到springcloudAlibaba博主用的是5.3.1的rocketmq引入封装好的消费者与生产者依赖用2.3.6就可以兼容使用。1、依赖引入!--rocketmq的版本管理-- dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-spring-boot-starter/artifactId version2.3.6/version /dependency2、entitypackage com.cloud.zhenyu.entity; import lombok.Data; import java.io.Serializable; /** 测试消息实体 */ Data public class DemoMsg implements Serializable { private static final long serialVersionUID 1L; private Long id; private String content; private Long timestamp; private String tag; private String hashKey; }3、producerpackage com.cloud.zhenyu.producer; import com.cloud.zhenyu.entity.DemoMsg; import org.apache.rocketmq.client.producer.SendCallback; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; /** RocketMQ 消息生产者 */ Component public class DemoRocketProducer { Autowired private RocketMQTemplate rocketMQTemplate; // 主题名称统一常量 public static final String TOPIC_DEMO topic-demo; /** 同步发送可靠等待Broker返回结果业务常用 */ public void sendSyncMsg(DemoMsg msg) { MessageDemoMsg message MessageBuilder.withPayload(msg).build(); SendResult sendResult rocketMQTemplate.syncSend(TOPIC_DEMO, message); System.out.println(同步发送结果 sendResult); } /** 异步发送不阻塞主线程回调接收结果 */ public void sendAsyncMsg(DemoMsg msg) { MessageDemoMsg message MessageBuilder.withPayload(msg).build(); rocketMQTemplate.asyncSend(TOPIC_DEMO, message, new SendCallback() { Override public void onSuccess(SendResult sendResult) { System.out.println(异步发送成功 sendResult); } Override public void onException(Throwable e) { System.err.println(异步发送失败 e.getMessage()); // 此处可做重试、落库补偿逻辑 } }); } /** 单向发送不等待响应日志类、不在乎是否送达场景 */ public void sendOneWayMsg(DemoMsg msg) { MessageDemoMsg message MessageBuilder.withPayload(msg).build(); rocketMQTemplate.sendOneWay(TOPIC_DEMO, message); } /** 带Tag过滤发送消费者可根据tag过滤消息 */ public void sendMsgWithTag(DemoMsg msg) { // topic:tag 格式 String destination TOPIC_DEMO : msg.getTag(); rocketMQTemplate.convertAndSend(destination, msg); } /** 顺序消息同一个hashKey的消息会进入同一队列保证有序 */ public void sendOrderMsg(DemoMsg msg) { rocketMQTemplate.syncSendOrderly(TOPIC_DEMO, msg, msg.getHashKey()); } }4、consumerpackage com.cloud.zhenyu.consumer; import com.cloud.zhenyu.entity.DemoMsg; import com.cloud.zhenyu.producer.DemoRocketProducer; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Service; /** 普通消息消费者 consumerGroup消费者组同一topic不同业务必须不同组 topic监听主题 selectorExpressiontag过滤* 接收全部tag1||tag2 接收指定tag / Slf4j Service RocketMQMessageListener( consumerGroup demo-consumer-group, topic DemoRocketProducer.TOPIC_DEMO, selectorExpression ) public class DemoRocketConsumer implements RocketMQListenerDemoMsg { /** 收到消息自动执行抛出异常会自动重试 */ Override public void onMessage(DemoMsg message) { log.info(消费者收到消息{}, message); // 业务处理逻辑 // 若抛出异常RocketMQ会重试投递正常执行无异常则标记消费成功 } }5、配置文件rocketmq: # NameServer 地址多个用分号分隔 name-server: 127.0.0.1:9876 # 生产者配置 producer: # 生产者组同业务统一组名 group: demo-producer-group # 发送消息超时时间 send-message-timeout: 3000 # 消费者统一配置可在这里也可注解单独指定三、三、本地测试说明博主本来是在阿里云部署测试玩的但是Rocketmq初始内存就要8G因为囊中羞涩所以只有将启动参数改到512M但是后续测试发现用本地java连接云服务器的mq会出现消费不了的异常原因如下Broker 注册给 NameServer 的是服务器内网 IP172.x.x.xWindows 本地无法访问内网网段。生产者仅连 9876 可发消息消费者拉取消息需连接 Broker 10911 内网地址直接失败表现发送成功、消费无响应 / 启动抛连接 null 异常。所以诸多麻烦下建议自己玩玩的话还是用windows本机把。四、消费者幂等与并发安全教学案例1、案例简介1.1 业务背景基于水务账单支付业务用户可在线缴纳水费账单支付完成后需完成三个核心操作1.生成缴费流水记录2.更新账单状态为【已缴费】3.多余支付金额自动充值至用户账户余额业务特性账单为系统预生成每期账单自动生成未支付时已存在唯一账单号支持多端登录、多人同账号操作、MQ异步消费存在极高并发数据风险。1.2 线上高频BUG场景1.同一账单多人并发重复支付多人登录同一账号同时支付同一未缴费账单生成不同支付单号的MQ消息导致重复扣款、重复生成流水。2.用户余额并发覆盖错乱同一用户同时触发缴费充值、退款扣费等操作代码查询余额后修改出现脏写覆盖资金对账错乱。2、解决方案采用「Redis分布式锁前置拦截 数据库原子行锁终极兜底」双层架构兼顾性能与绝对数据安全适用于所有支付、充值、退费等资金类MQ消费场景。第一层Redis分布式锁性能优化拦截99%并发•锁维度以账单唯一号 billNo 加锁核心业务主体•锁逻辑同一账单同一时间仅允许一条MQ消息执行业务•作用拦截瞬时并发重复请求避免大量无效数据库事务提升消费性能第二层数据库原子行锁终极兜底杜绝数据BUG解决Redis宕机、锁超时、主从切换等极端锁失效场景从数据库层面保证数据一致性包含两类原子锁1.账单状态行锁仅未支付账单可更新状态杜绝重复支付2.余额原子加减锁不查询余额直接数据库字段运算彻底解决并发覆盖错乱3、完整执行流程1.消费者监听支付MQ消息解析核心参数账单号billNo、用户ID、多余充值金额2.根据billNo尝试获取Redis分布式锁加锁失败则重试消息终止当前流程3.加锁成功开启数据库本地事务4.执行账单原子状态更新仅未支付账单更新成功已支付账单直接终止流程5.账单更新成功生成缴费流水记录6.执行用户余额原子加减操作处理多余充值金额7.事务整体提交业务执行成功手动ACK确认消息8.捕获异常实现重试机制超限消息转入死信队列9.finally释放分布式锁避免死锁4、代码实现import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQPushConsumerLifecycleListener; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Component; import org.springframework.transaction.support.TransactionTemplate; import com.alibaba.fastjson.JSON; import lombok.extern.slf4j.Slf4j; import javax.annotation.Resource; import java.nio.charset.StandardCharsets; import java.util.List; import java.util.concurrent.TimeUnit; /** 账单支付MQ消费者 解决多人并发重复支付、MQ重复消费、用户余额并发错乱 架构Redis账单锁前置拦截 数据库原子行锁兜底 / Slf4j Component RocketMQMessageListener( topic topic_pay_notify, consumerGroup pay-consumer-group, selectorExpression ) public class PayConsumer implements MessageListenerConcurrently, RocketMQPushConsumerLifecycleListener { private DefaultMQPushConsumer consumer; // 分布式锁过期时间大于业务最大执行时长防止业务未完成锁过期 private static final long LOCK_EXPIRE_TIME 30; Resource private RedisTemplateString, String redisTemplate; Resource private TransactionTemplate transactionTemplate; Resource private PayFlowService flowService; Resource private BillService billService; Resource private UserBalanceService balanceService; /** 初始化开启手动ACK */ Override public void prepareStart(DefaultMQPushConsumer consumer) { this.consumer consumer; // 关闭自动提交业务成功后手动ACK防止丢消息、重复消费异常 consumer.setManualCommit(true); } Override public ConsumeConcurrentlyStatus consumeMessage(ListMessageExt list, ConsumeConcurrentlyContext context) { MessageExt messageExt list.get(0); String billNo null; try { // 解析MQ消息 String body new String(messageExt.getBody(), StandardCharsets.UTF_8); PayDTO payDTO JSON.parseObject(body, PayDTO.class); billNo payDTO.getBillNo(); String lockKey pay:bill:lock: billNo; // 1.前置拦截账单维度分布式锁防止同一账单并发支付 boolean lockSuccess redisTemplate.opsForValue() .setIfAbsent(lockKey, PROCESSING, LOCK_EXPIRE_TIME, TimeUnit.SECONDS); if (!lockSuccess) { log.info(账单并发处理中稍后重试billNo{}, billNo); return ConsumeConcurrentlyStatus.RECONSUME_LATER; } // 2.事务执行全套支付业务 Boolean processSuccess transactionTemplate.execute(status -gt; { // 2.1 数据库原子更新账单状态仅未支付账单可更新行锁防重复支付 int updateRows billService.updateUnpaidBillStatus(billNo); if (updateRows lt; 0) { // 账单已支付无需重复处理 return false; } // 2.2 生成缴费流水记录 flowService.createPayFlow(payDTO); // 2.3 原子更新用户余额彻底解决并发错乱 if (payDTO.getSurplusAmount() gt; 0) { balanceService.addUserBalance(payDTO.getUserId(), payDTO.getSurplusAmount()); } return true; }); // 3.根据业务结果确认消息 if (processSuccess) { consumer.ack(messageExt); log.info(账单支付消费成功billNo{}, billNo); } else { consumer.ack(messageExt); log.warn(账单已支付过滤重复MQ消息billNo{}, billNo); } } catch (Exception e) { // 异常重试死信队列兜底 int retryTimes messageExt.getReconsumeTimes(); log.error(账单消费异常billNo{}重试次数{}, billNo, retryTimes, e); if (retryTimes gt; 3) { // 重试3次失败转入死信队列避免阻塞正常队列 consumer.sendToDeadLetter(messageExt); } else { return ConsumeConcurrentlyStatus.RECONSUME_LATER; } } finally { // 释放分布式锁防止死锁 if (billNo ! null) { redisTemplate.delete(pay:bill:lock: billNo); } } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }5、核心原子SQL### 防重复支付账单状态原子更新 UPDATE t_bill SET pay_status 1, pay_time NOW() WHERE bill_no #{billNo} AND pay_status 0; 防余额错乱用户余额原子加减直接字段运算规避查询赋值导致的并发覆盖问题支持缴费、退款、充值并发操作 -- 余额增加多余缴费、账户充值场景 UPDATE t_user SET balance balance #{amount} WHERE user_id #{userId}; -- 余额扣减水费扣费、退款场景 UPDATE t_user SET balance balance - #{amount} WHERE user_id #{userId} AND balance #{amount};6、总结1. 幂等维度选型核心资金、订单、账单类业务禁止用动态生成的流水号payNo做幂等必须用业务固有唯一主键billNo/orderNo。2. 并发分层思想Redis锁负责高性能前置拦截数据库原子锁负责最终数据兜底二者缺一不可。3. 余额并发铁律所有金额增减操作禁止查询后赋值必须使用数据库原子加减SQL。4. MQ可靠性规范核心资金业务必须开启手动ACK业务成功再提交异常自动重试超限转入死信队列。5. 杜绝伪幂等放弃「先查询后写入」的写法无法抵御高并发竞态问题。

相关新闻

MTAN在Visual Decathlon挑战赛中的应用:10个视觉任务同时学习实战指南

MTAN在Visual Decathlon挑战赛中的应用:10个视觉任务同时学习实战指南

MTAN在Visual Decathlon挑战赛中的应用:10个视觉任务同时学习实战指南 【免费下载链接】mtan The implementation of "End-to-End Multi-Task Learning with Attention" [CVPR 2019]. 项目地址: https://gitcode.com/gh_mirrors/mta/mtan 在计算机…

2026/7/21 5:57:10阅读更多 →
放置游戏设计学习指南:4 本必读书、GDC 演讲与 30 天实践路线

放置游戏设计学习指南:4 本必读书、GDC 演讲与 30 天实践路线

本文是「放置游戏策划实战」系列第 9 篇。 前八篇讨论了数值飞轮、独立乘区、蒙特卡洛模拟、系统循环、飞升系统、数据闭环与竞品拆解——重点都是“怎么做”。 这一篇换个角度:当你遇到具体设计问题时,应该查哪本书、看哪场演讲、看完马上做什么&#x…

2026/7/22 0:53:51阅读更多 →
最新360手机刷机 360手机官方刷机教程100%成功 360奇酷手机线刷救砖教程

最新360手机刷机 360手机官方刷机教程100%成功 360奇酷手机线刷救砖教程

360手机如何刷机 360手机官方刷机教程 360 手机刷机、解锁 BL、获取 ROOT 权限 视频教程 360手机刷机软件常见问题 问:360 手机早已停产,还支持刷机操作吗? 答:维护版正常刷机、解锁 Bootloader 以及获取 ROOT 权限 360手机&a…

2026/7/21 19:44:44阅读更多 →
Python字符串拼接性能优化:从+、join到f-string的实战指南

Python字符串拼接性能优化:从+、join到f-string的实战指南

1. 项目概述:为什么字符串拼接值得深究?刚接触Python那会儿,我也觉得字符串拼接不就是加号连一连的事儿吗?直到后来在项目中处理日志、拼接SQL、生成动态配置,甚至是在做性能敏感的数据处理时,才被现实狠狠…

2026/7/22 4:20:25阅读更多 →
Unity与React深度集成:WebView双向通信架构与工程实践

Unity与React深度集成:WebView双向通信架构与工程实践

1. 项目概述:为什么要在Unity WebView里跑React?这个问题乍一听有点“跨界”,一个做游戏和实时3D渲染的引擎,一个做现代Web前端界面的框架,它们俩怎么扯上关系了?但如果你正在开发一个需要复杂UI交互的3D应…

2026/7/22 4:20:25阅读更多 →
从面试翻车到生产落地,吃透长任务Agent的七大工程核心难点

从面试翻车到生产落地,吃透长任务Agent的七大工程核心难点

前段时间一位计算机硕士朋友面试头部AI基础设施公司算法岗,简历上亮眼的“企业级复杂流程Agent系统搭建”项目,本是他的加分王牌,结果被面试官几个直击生产痛点的问题问得全盘失语。 面试官没有追问复杂算法原理、模型微调技巧这些常规考点&a…

2026/7/22 4:20:25阅读更多 →
面试踩坑实录,为什么95%的程序员回答Agent意图识别,刚开口就失去录用机会

面试踩坑实录,为什么95%的程序员回答Agent意图识别,刚开口就失去录用机会

开篇:一场决定岗位去向的面试提问 近几年企业AI智能体岗位面试,有一道必考题反复出现在大厂、中腰部科技公司甚至传统行业数字化团队的笔面试环节,面试官轻描淡写抛出一句,就能快速筛掉绝大多数候选人,这个问题就是Age…

2026/7/22 4:20:25阅读更多 →
Windows LAPS本地管理员密码管理机制与部署实践

Windows LAPS本地管理员密码管理机制与部署实践

1. Windows LAPS 核心机制解析Windows LAPS(Local Administrator Password Solution)是微软针对企业环境中本地管理员账户密码管理难题设计的自动化解决方案。其核心工作原理可分解为三个关键环节:1.1 密码生命周期管理机制LAPS通过组策略客户…

2026/7/22 4:20:25阅读更多 →
PHP容器化部署与WebSocket服务在Kubernetes中的实践

PHP容器化部署与WebSocket服务在Kubernetes中的实践

1. 项目概述在云原生时代,PHP应用的容器化部署一直是个颇具挑战性的任务。最近我在阿里云ACK(Kubernetes)环境中成功部署了一个基于ThinkPHP框架的项目,并实现了WebSocket服务的搭建。整个过程踩了不少坑,也积累了一些…

2026/7/22 4:18:25阅读更多 →
Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/22 0:53:59阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/22 0:53:59阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/22 0:53:59阅读更多 →
中小企业小程序开发公司怎么选:预算、上手和售后避坑指南

中小企业小程序开发公司怎么选:预算、上手和售后避坑指南

中小企业做小程序,最常见的矛盾是预算有限,但又不希望功能太单薄;没有技术团队,但又希望后续能自己运营;想快速上线,又担心隐性收费和售后失联。选型时如果只看“低价套餐”或“案例数量”,很容…

2026/7/22 0:01:17阅读更多 →
GEO优化如何沉淀长期内容资产?广拓时代谈AI搜索时代的内容ROI

GEO优化如何沉淀长期内容资产?广拓时代谈AI搜索时代的内容ROI

企业做营销,最怕钱花完了,资产没有留下。 效果广告能带来一段时间的曝光,但预算停止后,流量往往也随之停止。短视频内容可能在几天内冲高,也可能很快沉下去。AI搜索时代,企业需要重新思考一个问题&#xff…

2026/7/22 0:01:17阅读更多 →
Agent 终态判定:何时该停止思考、给出最终回复

Agent 终态判定:何时该停止思考、给出最终回复

Agent 终态判定:何时该停止思考、给出最终回复 一、你的 Agent 在"再想想"的循环里绕了 12 轮,用户已经关窗口了 Agent 与人最大的区别是:人知道什么时候该停下来给答案,Agent 会一直"想"下去。你给 Agent 接…

2026/7/22 0:01:17阅读更多 →
YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

如果你在部署 YOLOv8 时,发现推理速度只有可怜的 1-2 FPS,而别人的演示视频却能跑到 30 FPS 以上,那么问题很可能不在模型本身,而在于你的整个处理链路。很多开发者拿到一个训练好的 YOLOv8 模型后,会直接使用官方示例…

2026/7/21 22:53:50阅读更多 →
Coze与Dify对比指南:低代码AI应用开发从入门到实战

Coze与Dify对比指南:低代码AI应用开发从入门到实战

1. 从零到一:为什么你需要了解 Coze 和 Dify?如果你对 AI 应用开发感兴趣,但一看到“大模型”、“智能体”、“工作流”这些词就头疼,觉得门槛太高,那这篇文章就是为你准备的。很多开发者,包括我自己&#…

2026/7/21 18:53:30阅读更多 →
AI生图工具怎么选?2026年6月版实测对比

AI生图工具怎么选?2026年6月版实测对比

做自媒体的朋友应该都有体会:配图一直是个让人头疼的问题。2026年,AI生图工具已经非常成熟了,但工具太多反而不知道怎么选。以下是截至2026年6月我对主流AI生图工具的实测对比。Midjourney V8.1:速度之王2026年6月11日&#xff0c…

2026/7/21 18:53:30阅读更多 →