SpringBoot 整合 RocketMq

1、安装 RocketMq

1docker search rocketmq 2 3docker pull rocketmqinc/rocketmq:4.4.0 4 5# 建立 name server 6docker run -d -p 9876:9876 -v /data/rocketmq/namesrv/logs:/root/logs -v /data/rocketmq/namesrv/store:/root/store --name rmqnamesrv -e "MAX_POSSIBLE_HEAP=100000000" rocketmqinc/rocketmq:4.4.0 sh mqnamesrv 7 8# 建立 broker.conf 9建立目录和broker配置文件: /data/rocketmq/conf/broker.conf 10 11brokerClusterName = DefaultCluster 12brokerName = broker-a 13brokerId = 0 14deleteWhen = 04 15fileReservedTime = 48 16brokerRole = ASYNC_MASTER 17flushDiskType = ASYNC_FLUSH 18brokerIP1 = 210.112.133.103 19 20# 建立 broker 21docker run -d -p 10911:10911 -p 10909:10909 -v /data/rocketmq/broker/logs:/root/logs -v /data/rocketmq/broker/store:/root/store -v /data/rocketmq/conf/broker.conf:/opt/rocketmq-4.4.0/conf/broker.conf --name rmqbroker --link rmqnamesrv:namesrv -e "NAMESRV_ADDR=namesrv:9876" -e "MAX_POSSIBLE_HEAP=200000000" rocketmqinc/rocketmq:4.4.0 sh mqbroker -c /opt/rocketmq-4.4.0/conf/broker.conf 22 23# 安装控制台 24docker pull styletang/rocketmq-console-ng 25docker run -e "JAVA_OPTS=-Drocketmq.config.namesrvAddr=210.112.133.103:9876 -Drocketmq.config.isVIPChannel=false" -p 8001:8080 -t styletang/rocketmq-console-ng 26

2、POM 加入 starter

1 <!-- https://mvnrepository.com/artifact/org.apache.rocketmq/rocketmq-spring-boot-starter --> 2 <dependency> 3 <groupId>org.apache.rocketmq</groupId> 4 <artifactId>rocketmq-spring-boot-starter</artifactId> 5 <version>2.1.1</version> 6 </dependency>

3、配置 rocketmq

1# rocketmq configuration 2rocketmq.name-server=210.112.133.103:9876 3rocketmq.producer.group=rocketmq_producer_group 4rocketmq.producer.send-message-timeout=3000

4、测试代码

1@ApiModel(description= "订单实体") 2@Data 3@NoArgsConstructor 4@AllArgsConstructor 5public class Order implements Serializable { 6 7 @ApiModelProperty(value = "主键") 8 private Long id; 9 @ApiModelProperty(value = "商品名称") 10 private String productName; 11 @ApiModelProperty(value = "价格") 12 private BigDecimal amount; 13 @ApiModelProperty(value = "时间") 14 private Date createTime; 15} 16 17 18/** 19 * 消息生产者 20 */ 21@Service 22public class RocketMqProducer { 23 24 private RocketMQTemplate rocketMQTemplate; 25 26 public RocketMqProducer(RocketMQTemplate rocketMQTemplate) { 27 this.rocketMQTemplate = rocketMQTemplate; 28 } 29 30 public RocketMQTemplate getRocketMQTemplate() { 31 return rocketMQTemplate; 32 } 33} 34 35 36@Api(value = "订单管理", tags = "管理订单增删改查") 37@Log4j2 38@RestController 39@RequestMapping("/order") 40public class OrderController { 41 42 private RocketMqProducer rocketMqProducer; 43 44 public OrderController(RocketMqProducer rocketMqProducer) { 45 this.rocketMqProducer = rocketMqProducer; 46 } 47 48 @ApiOperation(value = "保存或更新订单", notes = "id、productName、amount都是必输入项") 49 @ApiImplicitParams(value = { 50 @ApiImplicitParam(name = "id", value = "", dataTypeClass = Long.class, required = true), 51 @ApiImplicitParam(name = "productName", value = "", dataTypeClass = String.class, required = true), 52 @ApiImplicitParam(name = "amount", value = "", dataTypeClass = BigDecimal.class, required = true) 53 }) 54 @PostMapping("/saveOrUpdate") 55 public ResponseEntity<Result<Order>> saveOrUpdate( 56 @RequestParam Long id, 57 @RequestParam String productName, 58 @RequestParam BigDecimal amount 59 ) { 60 Order order = new Order(id, productName, amount, new Date()); 61 rocketMqProducer.getRocketMQTemplate().getProducer().setSendMsgTimeout(10000); 62 rocketMqProducer.getRocketMQTemplate().convertAndSend(Constants.ROCKET_MQ_ORDER, order); 63 return ResultUtils.ok(order); 64 } 65} 66 67@Log4j2 68@Service 69@RocketMQMessageListener(topic = Constants.ROCKET_MQ_ORDER, selectorExpression = "*", consumerGroup = Constants.ROCKET_MQ_CONSUMER_GROUP_ORDER) 70public class RocketMqOrderConsumerListener implements RocketMQListener<Order> { 71 @Override 72 public void onMessage(Order order) { 73 log.error("【onMessage start】"); 74 log.error(order.getId()); 75 log.error(order.getProductName()); 76 log.error(order.getAmount()); 77 log.error(order.getCreateTime()); 78 log.error("【onMessage end】"); 79 } 80}

5、控制台

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang