SpringBoot整合ActiveMQ消息队列(JmsTemplate实现方式)

整合 ROCKETMQ 实现消息生产消费(ROCKETMQTEMPLATE实现)

消息队列能解决什么问题:异步处理、应用解耦、流量削锋、日志处理、消息通讯

缺点:多出一个环节,需要保证消息队列的可用性。

MQTT

MQTT(Message Queuing Telemetry Transport,消息队列遥测传输)是IBM开发的一个即时通讯协议,ActiveMQ只是Apache下一个队列项目,不仅仅支持MQTT协议,也支持其他比如AMQP等协议。MQTT是协议,协议只是定义好的规则。ActiveMQ只是实现了MQTT协议的一个程序。

如果说传统的消息队列中间件一般应用于微服务之间,那么适用于物联网的微消息队列 MQTT 版则实现了端与云之间的消息传递和真正意义上的万物互联。

百度百科 AMQP https://baike.baidu.com/item/AMQP/8354716?fr=aladdin

ActiveMQ

ActiveMQ是一种开源的基于JMS(Java Message Servie)规范的一种消息中间件的实现,ActiveMQ的设计目标是提供标准的,面向消息的,能够跨越多语言和多系统的应用集成消息通信中间件。

5种消息类型

TextMessage:java.lang.String对象,如xml文件内容。 MapMessage:key/value键值对的集合,key是String对象,值类型可以是Java任何基本类型。 BytesMessage:字节流。 StreamMessage:Java 中的输入输出流。 ObjectMessage:Java中的可序列化对象。

SpringBoot实现JmsTemplate处理ActiveMQ消息

@EnableJms //启动消息队列
@EnableAutoConfiguration(exclude = DataSourceAutoConfiguration.class)
@SpringBootApplication
public class ActiveApplication {
          
   
    public static void main(String[] args) {
          
   
        SpringApplication.run(ActiveApplication.class, args);
    }
}

BeanConfig.java 定义存放消息的队列。

ProviderService 生产者,创建Controller测试

@Component
public class ProviderService {

    @Autowired private JmsMessagingTemplate jmsMessagingTemplate;

    @Resource(name = "deviceMSG")
    @Autowired private Queue queue;

    public void sendMessage(){
        Map<String, byte[]> msgMap = new HashMap<>();
        byte[] ctx = new byte[]{49, 54, 50, 47, 51, 48, 54, 51, 52, 50};
        String devId = "2201000001";
        msgMap.put("devId", devId.getBytes());
        msgMap.put("ctx", ctx);
        jmsMessagingTemplate.convertAndSend(queue ,msgMap);
    }
}

创建消费者 ConsumerService.java

@Component
public class ConsumerService {
          
   

    @Autowired private MainService mainService;

    @JmsListener(destination = "deviceMSG", concurrency = "10-20")
    public void handlerMessage(Map<String, byte[]> msgMap)
    {
          
   
        mainService.anylysisDataPackage(msgMap);
    }
}

创建Controller测试

@RestController
public class IndexController {
          
   
    @Autowired ProviderService providerService;
    
    @RequestMapping("/sendMessage")
    public void sendMessage(){
          
   
        logger.info("Provider send msg success");
        providerService.sendMessage();
    }
}

参考资料

mq原理&最佳实践 https://www.jianshu.com/p/2838890f3284

RabbitMQ MQTT协议和AMQP协议 https://www.cnblogs.com/bclshuai/p/8607517.html

分布式开放消息系统(RocketMQ)的原理与实践 http://www.jianshu.com/p/453c6e7ff81c

RocketMQ事务消费和顺序消费详解 http://www.cnblogs.com/520playboy/p/6750023.html

百度百科 AMQP https://baike.baidu.com/item/AMQP/8354716?fr=aladdin

经验分享 程序员 微信小程序 职场和发展