ActiveMQ

#1 前言 前一篇介绍了JMS有两种通信模型,一种是点对点通信,另一种是发布/订阅模型,本篇将会继续探讨这两种模型。本篇文章需要按照严谨的实验顺序才能获得相同的结果,这是因为消息持久化和持久订阅这两个特性的原因,在文章结尾和下一篇文章会做解答。** 所有的实验在启动之前都必须到管理后台删除相关的队列或者topic,否则数据也可能不同 **

ActiveMQ-5.11.x版本单机安装 中,由于设置了JMS通信端口的密码,所以下文在建立连接工厂时,都使用了密码。

为了方便,提取了一个专门创建connection的管理类。

1public class ActiveMQManager { 2 private static ConnectionFactory connectionFactory; 3 static { 4 // connectionFactory = new ActiveMQConnectionFactory("admin", "123456", "tcp://192.168.88.18:61616"); 5 connectionFactory = new ActiveMQConnectionFactory("tcp://127.0.0.1:61616"); 6 } 7 public static Connection createConnection() throws JMSException { 8 return connectionFactory.createConnection(); 9 } 10}

OK,下面开始正题吧。

#2 点对点模型 本节将实现消费者通过队列"test-queue"给消费者发布消息,消费者以异步方式获取消息并打印消息。先把消息生产者和消费者的代码贴一下。

##2.1 实验代码 ** 消息生产者 **:

1public class Producer { 2 public static final String QUEUE_NAME = "test-queue"; 3 4 public static void main(String[] args) { 5 System.out.println("Producer started!"); 6 7 String message_body = "消息 : " + System.currentTimeMillis(); 8 9 try { 10 //获取连接 11 Connection connection = ActiveMQManager.createConnection(); 12 13 //启动连接 14 connection.start(); 15 16 //开启会话,第一个参数指定是否使用事务,第二个参数指示消费者是否需要手动应答自己已经接收到消息 17 Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE); 18 19 //建立队列 20 Queue queue = session.createQueue(QUEUE_NAME); 21 22 //获取消息生产者对象 23 MessageProducer producer = session.createProducer(queue); 24 25 //建立消息对象 26 Message message = session.createTextMessage(message_body); 27 28 //发送消息 29 producer.send(message); 30 31 System.out.println("成功发送消息:" + message_body); 32 33 producer.close(); 34 session.close(); 35 connection.close(); 36 } catch (Exception e) { 37 e.printStackTrace(); 38 } 39 40 System.out.println("Producer end!"); 41 } 42}

** 消息消费者 **:

1public class Consumer { 2 public static void main(String[] args) throws IOException { 3 System.out.println("Consumer started!"); 4 5 try { 6 //获取连接 7 Connection connection = ActiveMQManager.createConnection(); 8 9 //启动连接 10 connection.start(); 11 12 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); 13 14 //建立队列,与消息生产者使用的队列名一致 15 Queue queue = session.createQueue(Producer.QUEUE_NAME); 16 17 //建立消费者 18 MessageConsumer consumer = session.createConsumer(queue); 19 20 //监听消息 21 consumer.setMessageListener(new MessageListener() { 22 @Override 23 public void onMessage(Message message) { 24 try { 25 TextMessage textMessage = (TextMessage) message; 26 String text = textMessage.getText(); 27 System.out.println("Consumer 获取消息 ---->" + text); 28 } catch (Exception e) { 29 e.printStackTrace(); 30 } 31 32 } 33 }); 34 35 // 线程一直等待 36 System.in.read(); 37 38 consumer.close(); 39 session.close(); 40 connection.stop(); 41 42 } catch (Exception e) { 43 e.printStackTrace(); 44 } 45 46 System.out.println("Consumer end!"); 47 } 48}

2.2 启动程序和小结论

** (1)顺序:启动MQ broker服务 -> 启动消息发布者 -> 启动消息消费者 **
消息可以被消费者正常接收。

** (2)顺序:启动MQ broker服务 -> 启动消息发布者 -> 停止MQ broker服务 -> 启动MQ broker服务 -> 启动消息消费者 ** 消息可以被消费者正常接收,通过查看消息内容,可以知道该消息是停止MQ服务之前发送的。

** (3)顺序:启动MQ broker服务 -> 启动消息消费者 -> 启动消息发布者 ** 消息可以被消费者正常接收。

看到这里,大家可能会感到奇怪,为什么MQ服务重启后,之前发布的消息仍然可以被后面启动的消费者收到呢?这与activeMQ的持久化消息机制有关。下一篇会对此做解答。

#3 发布/订阅模型 本节将定义两个消息订阅者,它们共同订阅名叫"test-topic"这个主题,消息生产者发布消息后,两个消费者都能接收到同一个消息。

##2.1 实验代码 下面是生产者和订阅者的代码。

** 消息生产者 **

1public class Producer { 2 public static final String TOPIC_NAME = "test-topic"; 3 4 public static void main(String[] args) { 5 System.out.println("Producer started!"); 6 7 String message_body = "消息 : " + System.currentTimeMillis(); 8 9 try { 10 //获取连接 11 Connection connection = ActiveMQManager.createConnection(); 12 13 //启动连接 14 connection.start(); 15 16 //开启会话,第一个参数指定是否使用事务,第二个参数指示消费者是否需要手动应答自己已经接收到消息 17 Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE); 18 19 //建立话题 20 Topic topic = session.createTopic(TOPIC_NAME); 21 22 //获取消息生产者对象 23 MessageProducer producer = session.createProducer(topic); 24 25 //建立消息对象 26 Message message = session.createTextMessage(message_body); 27 28 //发送消息 29 producer.send(message); 30 31 System.out.println("成功发送消息:" + message_body); 32 33 producer.close(); 34 session.close(); 35 connection.close(); 36 } catch (Exception e) { 37 e.printStackTrace(); 38 } 39 40 System.out.println("Producer end!"); 41 } 42}

** 消息订阅者 ** :

1public class Subscriber1 { 2 public static void main(String[] args) throws IOException { 3 System.out.println("Subscriber1 started!"); 4 5 try { 6 //获取连接 7 Connection connection = ActiveMQManager.createConnection(); 8 9 //启动连接 10 connection.start(); 11 12 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); 13 14 //建立话题 15 Topic topic = session.createTopic(Producer.TOPIC_NAME); 16 17 //建立消费者 18 MessageConsumer consumer = session.createConsumer(topic); 19 20 //监听消息 21 consumer.setMessageListener(new MessageListener() { 22 @Override 23 public void onMessage(Message message) { 24 try { 25 TextMessage textMessage = (TextMessage) message; 26 String text = textMessage.getText(); 27 System.out.println("Subscriber1 获取消息 ---->" + text); 28 } catch (Exception e) { 29 e.printStackTrace(); 30 } 31 32 } 33 }); 34 35 // 线程一直等待 36 System.in.read(); 37 38 consumer.close(); 39 session.close(); 40 connection.stop(); 41 42 } catch (Exception e) { 43 e.printStackTrace(); 44 } 45 46 System.out.println("Subscriber1 end!"); 47 } 48}

3.2 启动程序和小结论

** (1)顺序:启动MQ broker服务 -> 启动消息发布者 -> 启动两个消息订阅者 **
消息订阅者没有接收到任何消息。

** (2)顺序:启动MQ broker服务 -> 启动两个消息订阅者 -> 启动消息发布者 **
两个消息订阅者都可以被消费者正常接收。

** (3)顺序:启动MQ broker服务 -> 启动消息发布者 -> 停止MQ broker服务 -> 启动MQ broker服务 -> 启动两个消息订阅者 **

wtf!消息没有接收到,如果是重要的消息,那岂不是要哭死了。

对于第三条,由于MQ服务重启导致订阅者消息获取不到,activeMQ提供了解决方法,即持久主题订阅(durable topic subscription)机制,嗯,下一篇将会对该机制详细解读。

代码

http://git.oschina.net/thinwonton/activemq-showcase

点赞
收藏

评论区

加载中...

相关推荐

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 )