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、控制台
