ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

spring cloud alibaba2022版本集成RocketMQ

spring cloud alibaba2022版本集成RocketMQ 目录一、安装部署RocketMQ二、项目添加依赖包2.1、pom.xml引入RocketMQ依赖包2.2、配置新增三、项目编写业务代码3.1、新增订单和取消订单方法做消息通知四、测试消息发送与消费4.1、顺序消费4.2、普通消费spring cloud alibaba2022版本分布式框架搭建示例之集成RocketMQ。一、安装部署RocketMQ1.1、下载包安装地址下载 | RocketMQ我这里下载了5.5.0版本解压后配置ROCKET_HOME环境变量,指定rocketmq的安装包文件路径。在rocketMQ安装包的bin目录下找到 runserver.cmd 文件右键「编辑」用记事本 /或Notepad 打开RocketMQ 默认配置内存要求比较高需先修改启动脚本以便可以在个人电脑上跑起来.。上图中的两个地方改成如下set JAVA_OPT%JAVA_OPT% -server -Xms512m -Xmx512m -Xmn512m -XX:MetaspaceSize128m -XX:MaxMetaspaceSize320m找到mqnamesrv.cmd文件双击运行。新起一个cmd窗口在安装包的bin目录下执行命令关联 NameServer 允许自动创建 Topicmqbroker -n 127.0.0.1:9876 autoCreateTopicEnabletrue启动成功如下二、项目添加依赖包补充——父项目pom.xml?xml version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdorg.example/groupId artifactIdspring-cloud-alibaba/artifactId version1.0-SNAPSHOT/version packagingpom/packaging modules modulegoods-service/module moduleorder-service/module modulegateway-service/module /modules properties maven.compiler.source17/maven.compiler.source maven.compiler.target17/maven.compiler.target project.build.sourceEncodingUTF-8/project.build.sourceEncoding spring-boot.version3.0.13/spring-boot.version spring-cloud.version2022.0.2/spring-cloud.version spring-cloud-starter-bootstrap.version3.1.5/spring-cloud-starter-bootstrap.version spring-cloud-alibaba.version2022.0.0.0/spring-cloud-alibaba.version /properties dependencyManagement dependencies !-- Spring Boot 3.x -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version${spring-boot.version}/version typepom/type scopeimport/scope /dependency !-- Spring Cloud 2022.x -- dependency groupIdorg.springframework.cloud/groupId artifactIdspring-cloud-dependencies/artifactId version${spring-cloud.version}/version typepom/type scopeimport/scope /dependency !-- Spring Cloud Bootstrap 3.x -- dependency groupIdorg.springframework.cloud/groupId artifactIdspring-cloud-starter-bootstrap/artifactId version${spring-cloud-starter-bootstrap.version}/version /dependency !-- Alibaba Cloud 2.x -- dependency groupIdcom.alibaba.cloud/groupId artifactIdspring-cloud-alibaba-dependencies/artifactId version${spring-cloud-alibaba.version}/version typepom/type scopeimport/scope /dependency !-- Nacos注册中心 -- dependency groupIdcom.alibaba.cloud/groupId artifactIdspring-cloud-starter-alibaba-nacos-discovery/artifactId version${spring-cloud-alibaba.version}/version /dependency !-- Nacos配置中心 -- dependency groupIdcom.alibaba.cloud/groupId artifactIdspring-cloud-starter-alibaba-nacos-config/artifactId version${spring-cloud-alibaba.version}/version /dependency !-- Spring Boot 3.x -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-actuator/artifactId version${spring-boot.version}/version /dependency !-- Spring Boot 3.x -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId version${spring-boot.version}/version /dependency dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId version1.18.30/version /dependency !-- RocketMQ -- dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-spring-boot-starter/artifactId version2.3.4/version scopecompile/scope exclusions !-- 排出低版本冲突 -- exclusion groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId /exclusion /exclusions /dependency dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId version5.3.1/version /dependency /dependencies /dependencyManagement /project2.1、pom.xml引入RocketMQ依赖包SpringCloud Alibaba2022版本对应rocketMQ版本如下dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-spring-boot-starter/artifactId version2.3.4/version scopecompile/scope exclusions !-- 排出低版本冲突 -- exclusion groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId /exclusion /exclusions /dependency dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId version5.3.1/version /dependency2.2、配置新增rocketmq: name-server: 127.0.0.1:9876 # 生产者配置 producer: group: order-group retry-times-when-send-failed: 3 send-message-timeout: 3000 compress-message-body-threshold: 4096 # 消息压缩阈值4KB # 消费者配置 consumer: group: order-group consume-thread-max: 20 # 最大消费线程数三、项目编写业务代码3.1、新增订单和取消订单方法做消息通知消息发送方package org.example.service.impl; import com.alibaba.fastjson.JSON; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import net.minidev.json.JSONObject; import org.apache.rocketmq.client.producer.SendCallback; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.example.controller.req.OrderRequest; import org.example.entity.Order; import org.example.feign.GoodsFeignClient; import org.example.feign.res.GoodsResponse; import org.example.mapper.OrderMapper; import org.example.service.IOrderService; import org.example.service.impl.dto.GoodsDTO; import org.springframework.stereotype.Service; import java.math.BigDecimal; import java.text.SimpleDateFormat; import java.time.LocalDateTime; import java.util.Date; import java.util.Objects; import java.util.Random; /** * p * 服务实现类 * /p * * author ChengJiangBo * since 2026-07-02 */ Slf4j Service public class OrderServiceImpl extends ServiceImplOrderMapper, Order implements IOrderService { Resource GoodsFeignClient goodsFeignClient; Resource RocketMQTemplate rocketMQTemplate; Override public boolean save(OrderRequest orderRequest) { Order order new Order(); if(Objects.nonNull(orderRequest.getGoodsId())){ GoodsResponse goodsResponse goodsFeignClient.getGoodsById(orderRequest.getGoodsId()); if(Objects.nonNull(goodsResponse)){ if(goodsResponse.getPrice().compareTo(new BigDecimal(0)) 0){ log.info(未查询到商品下单失败); return false; } order.setGoodsPrice(goodsResponse.getPrice()); order.setGoodsNum(orderRequest.getGoodsNum()); order.setGoodsVersion(goodsResponse.getGoodsVersion()); order.setAmount(new BigDecimal(orderRequest.getGoodsNum()).multiply(goodsResponse.getPrice())); order.setGoodsId(goodsResponse.getId()); order.setGoodsVersion(goodsResponse.getGoodsVersion()); SimpleDateFormat sdf new SimpleDateFormat(yyyyMMddHHmmss); order.setOrderNo( sdf.format(new Date()) new Random().nextInt(10000000) ); order.setCreateTime(LocalDateTime.now()); order.setUpdateTime(LocalDateTime.now()); boolean result super.save(order); log.info(下单成功订单编号{}, order.getOrderNo()); if (result) { GoodsDTO goodsDTO new GoodsDTO(); goodsDTO.setGoodsId(order.getGoodsId()); goodsDTO.setGoodsNum(order.getGoodsNum()); goodsDTO.setOrderNo(order.getOrderNo()); rocketMQTemplate.asyncSend(goods-topic:inventory-update, order, new SendCallback() { Override public void onSuccess(SendResult sendResult) { log.info(异步发送成功发送结果{}, sendResult); } Override public void onException(Throwable e) { log.info(异步发送失败发送结果{}, e); } }); } return result; } } return false; } Override public boolean cancel(String orderNo) { log.info(取消订单订单编号{}, orderNo); Order order getOne(new QueryWrapperOrder().eq(order_no, orderNo)); if(Objects.nonNull(order)){ GoodsResponse goods goodsFeignClient.getGoodsById(order.getGoodsId()); if(Objects.nonNull(goods)){ rocketMQTemplate.asyncSend(order-topic:order_cancel, JSON.toJSONString(goods), new SendCallback() { Override public void onSuccess(SendResult sendResult) { log.info(异步发送成功发送结果{}, sendResult); order.setOrderStatus(2); // Set order status to canceled updateById(order); } Override public void onException(Throwable e) { log.info(异步发送失败发送结果{}, e); } }); } } return true; } }package org.example.service.impl.dto; import lombok.Getter; import lombok.Setter; Setter Getter public class GoodsDTO{ private Long goodsId; private Integer goodsNum; private String orderNo; Override public String toString() { return GoodsDTO{ goodsId goodsId , goodsNum goodsNum , orderNo orderNo \ }; } }消息接收方rocketmq: name-server: 127.0.0.1:9876 # 生产者配置 producer: group: goods-group retry-times-when-send-failed: 3 send-message-timeout: 3000 compress-message-body-threshold: 4096 # 消息压缩阈值4KB # 消费者配置 consumer: group: goods-group consume-thread-max: 20 # 最大消费线程数新增两个消费者如下package org.example.application; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly; import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.spring.annotation.ConsumeMode; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.apache.rocketmq.spring.support.RocketMQConsumerLifecycleListener; import org.example.application.dto.GoodsDTO; import org.springframework.stereotype.Component; /** * 商品模块的MQ消费者 */ Slf4j Component RocketMQMessageListener(topic goods-topic, consumerGroup test-group, consumeMode ConsumeMode.ORDERLY ) public class GoodsMQConsumer implements RocketMQListenerGoodsDTO, RocketMQConsumerLifecycleListenerDefaultMQPushConsumer { Override public void onMessage(GoodsDTO message) { log.info(顺序消费消息{}, message); // 手动处理消息需确保幂等性 } Override public void prepareStart(DefaultMQPushConsumer consumer) { //手动配置ACK(默认是自动ACK) consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET); // 从最后一条消息开始消费 consumer.registerMessageListener((MessageListenerOrderly) (msgs, context)-{ // 顺序消费的手动ACK逻辑 context.setAutoCommit(true); // 自动提交顺序消费建议开启 return ConsumeOrderlyStatus.SUCCESS; }); } }package org.example.application; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.ConsumeMode; import org.apache.rocketmq.spring.annotation.MessageModel; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; /** * 商品模块的MQ消费者 */ Slf4j Component RocketMQMessageListener( topic order-topic, consumerGroup test-group, selectorExpression order_refund || order_cancel, // 过滤多个Tag consumeMode ConsumeMode.CONCURRENTLY, // 并发消费默认 messageModel MessageModel.CLUSTERING // 集群消费默认 ) public class OrderMQConsumer implements RocketMQListenerString { Override public void onMessage(String message) { log.info(消费消息{}, message); // 手动处理消息需确保幂等性 } }四、测试消息发送与消费4.1、顺序消费apifox执行创建订单接口order-service控制台打印结果goods-service控制台打印结果4.2、普通消费执行取掉订单接口order-service控制台打印结果goods-service控制台打印结果
返回列表