20200202 ActiveMQ 3. Java编码实现ActiveMQ通讯

ActiveMQ 3. Java编码实现ActiveMQ通讯

3.1. 队列(Queue)

目的地(Destination)分为:

  • 点对点的队列(Queue)
  • 一对多的主题(Topic)

3.1.1. 上手代码

  1. pom.xml

    <dependency> <groupId>org.apache.activemq</groupId> <artifactId>activemq-all</artifactId> <version>5.15.9</version> </dependency> <dependency> <groupId>org.apache.xbean</groupId> <artifactId>xbean-spring</artifactId> <version>3.16</version> </dependency>
  2. 生产者代码

    public class JmsProducer {

    1public static final String ACTIVEMQ_URL = "tcp://192.168.181.128:61616/"; 2public static final String QUEUE_NAME = "queue01"; 3 4public static void main(String[] args) throws JMSException { 5 // 1. 创建连接工厂,按照给定的URL地址,采用默认用户名密码 6 ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL); 7 // 2. 通过连接工厂,获得连接Connection并启动访问 8 Connection connection = activeMQConnectionFactory.createConnection(); 9 connection.start(); 10 11 // 3. 创建会话Session 12 // 两个参数,第一个是事务控制,第二个是签收控制 13 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); 14 // 4. 创建目的地(具体是队列queue或主题topic) 15 Queue queue = session.createQueue(QUEUE_NAME); 16 // 5. 创建消息的生产者 17 MessageProducer messageProducer = session.createProducer(queue); 18 // 6. 通过消息生产者发送消息 19 for (int i = 0; i < 3; i++) { 20 // 7. 创建消息 21 TextMessage textMessage = session.createTextMessage("msg---" + i); 22 // 8. 发送给MQ 23 messageProducer.send(textMessage); 24 } 25 // 9. 关闭资源 26 messageProducer.close(); 27 session.close(); 28 connection.close(); 29 30 System.out.println("*****消息发布到MQ完成*****"); 31}

    }

  3. 消费者代码

    public class JmsConsumer {

    1public static final String ACTIVEMQ_URL = "tcp://192.168.181.128:61616/"; 2public static final String QUEUE_NAME = "queue01"; 3 4public static void main(String[] args) throws JMSException { 5 // 1. 创建连接工厂,按照给定的URL地址,采用默认用户名密码 6 ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL); 7 // 2. 通过连接工厂,获得连接Connection并启动访问 8 Connection connection = activeMQConnectionFactory.createConnection(); 9 connection.start(); 10 11 // 3. 创建会话Session 12 // 两个参数,第一个是事务控制,第二个是签收控制 13 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); 14 // 4. 创建目的地(具体是队列queue或主题topic) 15 Queue queue = session.createQueue(QUEUE_NAME); 16 17 // 5. 创建消费者 18 MessageConsumer messageConsumer = session.createConsumer(queue); 19 while (true) { 20 TextMessage textMessage = (TextMessage) messageConsumer.receive(); 21 if (textMessage != null) { 22 System.out.println("*****消费者收到消息:" + textMessage.getText()); 23 } else { 24 break; 25 } 26 } 27 28 // 6. 关闭资源 29 messageConsumer.close(); 30 session.close(); 31 connection.close(); 32}

    }

3.1.2. receive()方法说明

1// 收到消息前一直阻塞进程 2javax.jms.MessageConsumer#receive() 3// 超时后不再阻塞进程 4javax.jms.MessageConsumer#receive(long timeout)

3.1.3. 消费者监听器方式接收消息

监听器方式属于异步非阻塞方式,所以需要手动阻塞进程

1messageConsumer.setMessageListener(new MessageListener() { 2 @SneakyThrows 3 @Override 4 public void onMessage(Message message) { 5 if (null != message && message instanceof TextMessage) { 6 System.out.println("消费者监听器监听到消息***********" + ((TextMessage) message).getText()); 7 } 8 } 9}); 10// 手动阻塞进程 11System.in.read();

3.1.4. 消费者三大消费情况

  1. 先生产,只启动1号消费者。问题:1号消费者可以消费消息吗?

    可以

  2. 先生产,先启动1号消费者,再启动2号消费者。问题:2号消费者可以消费消息吗?

    1号消费者可以消费消息;2号消费者不可以消费消息;

  3. 先启动2个消费者,再生产6条消息。问题:消费情况如何?

    2个消费者各消费一半消息;

3.1.5. 两种消费方式

  1. 同步阻塞方式(receive()

  2. 异步非阻塞方式(消费者监听器onMessage()

3.1.6. 点对点消息传递域的特点

  1. 每个消息只能有一个消费者,类似1对1的关系,类似于快递

  2. 消息的消费者和生产者没有时间上的相关性,类似于短信

  3. 消息被消费后队列中不会再存储,所以消费者不会消费到已经被消费掉的消息

3.2. 主题(Topic)

3.2.1. 发布订阅消息传递域的特点

  1. 每个消息可以有多个消费者,属于一对多的关系
  2. 生产者和消费者有时间上的相关性,订阅一个主题的消费者只能消费自它订阅之后发布的消息
  3. 生产者生产时,topic不保存消息,它是无状态的不落地,假如无人订阅就去生产,那就是一条废消息,所以,一般先启动消费者再启动生产者

JMS规范允许客户创建持久订阅,这在一定程度上放松了时间上的相关性要求。持久订阅允许消费者消费它在未处于激活状态时的消息。一句话,类似微信公众号订阅

3.2.2. 上手代码

测试时要先启动消费者,后启动生产者。

  1. 生产者代码

    public class JmsProducer_Topic {

    1public static final String ACTIVEMQ_URL = "tcp://192.168.181.128:61616/"; 2public static final String TOPIC_NAME = "topic01"; 3 4public static void main(String[] args) throws JMSException { 5 // 1. 创建连接工厂,按照给定的URL地址,采用默认用户名密码 6 ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL); 7 // 2. 通过连接工厂,获得连接Connection并启动访问 8 Connection connection = activeMQConnectionFactory.createConnection(); 9 connection.start(); 10 11 // 3. 创建会话Session 12 // 两个参数,第一个是事务控制,第二个是签收控制 13 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); 14 // 4. 创建目的地(具体是队列queue或主题topic) 15 Topic topic = session.createTopic(TOPIC_NAME); 16 // 5. 创建消息的生产者 17 MessageProducer messageProducer = session.createProducer(topic); 18 // 6. 通过消息生产者发送消息 19 for (int i = 0; i < 3; i++) { 20 // 7. 创建消息 21 TextMessage textMessage = session.createTextMessage("topic---" + i); 22 // 8. 发送给MQ 23 messageProducer.send(textMessage); 24 } 25 // 9. 关闭资源 26 messageProducer.close(); 27 session.close(); 28 connection.close(); 29 30 System.out.println("*****topic消息发布到MQ完成*****"); 31}

    }

  2. 消费者代码

    public class JmsConsumer_Topic {

    1public static final String ACTIVEMQ_URL = "tcp://192.168.181.128:61616/"; 2public static final String TOPIC_NAME = "topic01"; 3 4public static void main(String[] args) throws JMSException, IOException { 5 System.out.println("我是1号消费者"); 6 // System.out.println("我是2号消费者"); 7 // System.out.println("我是3号消费者"); 8 9 // 1. 创建连接工厂,按照给定的URL地址,采用默认用户名密码 10 ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL); 11 // 2. 通过连接工厂,获得连接Connection并启动访问 12 Connection connection = activeMQConnectionFactory.createConnection(); 13 connection.start(); 14 15 // 3. 创建会话Session 16 // 两个参数,第一个是事务控制,第二个是签收控制 17 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); 18 // 4. 创建目的地(具体是队列queue或主题topic) 19 Topic topic = session.createTopic(TOPIC_NAME); 20 21 // 5. 创建消费者 22 MessageConsumer messageConsumer = session.createConsumer(topic); 23 24 25 messageConsumer.setMessageListener(new MessageListener() { 26 @SneakyThrows 27 @Override 28 public void onMessage(Message message) { 29 if (null != message && message instanceof TextMessage) { 30 System.out.println("消费者监听器监听到 TOPIC 消息***********" + ((TextMessage) message).getText()); 31 } 32 } 33 }); 34 // 手动阻塞进程 35 System.in.read(); 36 37 // 6. 关闭资源 38 messageConsumer.close(); 39 session.close(); 40 connection.close(); 41}

    }

3.3. 两种模式比较

比较项目

Topic 模式

Queue模式

工作模式

订阅-发布”模式,如果当前没有订阅者,消息将会被丢弃;如果有多个订阅者,那么这些订阅者都会收到消息

"负载均衡"模式,如果当前没有消费者,消息也不会丢弃;如果有多个消费者,那么一条消息只会发送给其中一个消费者,并且要求消费者ack消息

有无状态

无状态

Queue数据默认会在MQ服务器上以文件形式保存。也可以配置成DB存储

传递完整性

如果没有订阅者,消息会被丢弃

消息不会被丢弃

处理效率

由于消息要按照订阅者数量进行复制,所以处理性能会随着订阅者的增加而明显降低,并且还要结合不同消息协议自身的性能差异

由于一条消息只发送给一个消费者,所以就算消费者再多,性能也不会明显降低。当然不同消息协议的具体性能也是有差异的

点赞
收藏

评论区

加载中...

相关推荐

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

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )