ActiveMQ入门系列三:发布/订阅模式

在上一篇《ActiveMQ入门系列二:入门代码实例(点对点模式)》中提到了ActiveMQ中的两种模式:点对点模式(PTP)和发布/订阅模式(Pub & Sub),详细介绍了点对点模式并用代码实例进行说明,今天就介绍下发布/订阅模式。

一、理论基础

发布/订阅模式的工作示意图:

  • 消息生产者将消息(发布)到topic中,可以同时有多个消息消费者(订阅)消费该消息。
  • 和点对点方式不同,发布到topic的消息会被所有订阅者消费。
  • 当生产者发布消息,不管是否有消费者,都不会保存消息。
  • 一定要先有消息的消费者,后有消息的生产者。

二、代码实现

  1. 生产者

    package com.sam.topic;
    
    import org.apache.activemq.ActiveMQConnectionFactory;
    
    import javax.jms.*;
    
    /**
     * @author JAVA开发老菜鸟
     *
     */
    public class TopicProducer {
    
        public static final  String QUEUE_NAME = "topic-demo";//队列名
    
        public void producer(String message) throws JMSException {
            ConnectionFactory factory = null;
            Connection connection = null;
            Session session = null;
            MessageProducer producer = null;
            try {
                /**
                 * 1.创建连接工厂
                 * 创建工厂,构造方法有三个参数:分别是用户名、密码、连接地址
                 * 无参构造:有默认的连接地址,localhost
                 * 一个参数:无验证模式,无用户的认证
                 * 三个参数:有认证和连接地址
                 */
                factory = new ActiveMQConnectionFactory("admin","admin","tcp://localhost:61616");
                /**
                 * 2.创建连接
                 * 无参数
                 * 有参数:用户名、密码
                 */
                connection = factory.createConnection();
                /**
                 * 3.启动连接
                 * 生产者可以不启动,因为在发送消息的时候回进行检查
                 * 如果未启动连接,会自动启动
                 * 如果有特殊配置,需要配置完成后再启动连接
                 */
                connection.start();
                /**
                 * 4.用连接创建会话
                 * 有两个参数:是否需要事务、消息确认机制
                 * 如果支持事务,对于生产者来说第二个参数就无效了,建议传入Session.SESSION_TRANSACTED
                 * 如果不支持事务,第二个参数必须传递且有效
                 *
                 * AUTO_ACKNOWLEDGE:自动确认,消息处理后自动确认(商业开发不推荐)
                 * CLIENT_ACKNOWLEDGE:客户端手动确认,消费者处理后必须手动确认
                 * DUPS_OK_ACKNOWLEDGE:有副本的客户端手动确认,消息可以多次处理(不建议)
                 */
                session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
                /**
                 * 5.用会话创建目的地(主题)、生产者、消息
                 * 队列名是队列的唯一标记
                 * 创建生产者的时候可以不指定目的地,可以在发送的时候指定
                 */
                Destination destination = session.createTopic(QUEUE_NAME);
                producer = session.createProducer(destination);
                TextMessage textMessage = session.createTextMessage(message);
                /**
                 * 6.生产者发送消息到目的地
                 */
                producer.send(textMessage);
                System.out.println("消息发送成功");
            } catch(Exception ex){
                throw ex;
            } finally {
                /**
                 * 7.释放资源
                 */
                if(producer != null){
                    producer.close();
                }
    
                if(session != null){
                    session.close();
                }
    
                if(connection != null){
                    connection.close();
                }
            }
        }
    
        public static void main(String[] args){
            TopicProducer producer = new TopicProducer();
            try{
                producer.producer("hello, activemq");
            } catch (Exception ex){
                ex.printStackTrace();
            }
        }
    }

    发布/订阅模式的生产者和点对点模式的代码主要区别就是Destination的创建方式,点对点模式是调用session.createQueue(QUEUE_NAME),而发布/订阅模式是调用session.createTopic(QUEUE_NAME)。

  2. 消费者

    package com.sam.topic;
    
    import org.apache.activemq.ActiveMQConnectionFactory;
    
    import javax.jms.*;
    import java.io.IOException;
    
    /**
     * @author JAVA开发老菜鸟
     *
     * 观察者消费--监听消费
     */
    public class TopicConsumer {
    
        public void consumer() throws JMSException, IOException {
            ConnectionFactory factory = null;
            Connection connection = null;
            Session session = null;
            MessageConsumer consumer = null;
            try {
                factory = new ActiveMQConnectionFactory("admin","admin","tcp://localhost:61616");
                connection = factory.createConnection();
                /**
                 * 消费者必须启动连接,否则无法消费
                 */
                connection.start();
                session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
                Destination destination = session.createTopic(TopicProducer.QUEUE_NAME);
                consumer = session.createConsumer(destination);
                /**
                 * 注册监听器,队列中的消息变化会自动触发监听器,接收并自动处理消息
                 *
                 * 监听器一旦注册,永久有效,一直到程序关闭
                 * 监听器可以注册多个,相当于集群
                 * activemq自动轮询多个监听器,实现并行处理
                 */
                consumer.setMessageListener(new MessageListener() {
                    @Override
                    public void onMessage(Message message) {
    
                        try {
                            TextMessage om = (TextMessage) message;
                            String data = om.getText();
                            System.out.println(data);
                        } catch (JMSException e) {
                            e.printStackTrace();
                        }
                    }
                });
            } catch(Exception ex){
                throw ex;
            }
        }
    
        public static void main(String[] args){
            TopicConsumer consumer = new TopicConsumer();
            try{
                consumer.consumer();
            } catch (Exception ex){
                ex.printStackTrace();
            }
        }
    }

    消费者在点对点监听消费的基础上进行变化,主要区别有两个:1.同生产者一样,也是Destination的创建方式不同; 2.消息无需手动确认,直接采用自动确认机制

代码写完了,接下来进行测试,由于subscribe可以有多个,而且每个都可以消费到相同的消息,因此我们消费者启动两个。

先执行生产者

在控制台页面的Topics下出现了我定义的topic并且有1条消息发送成功且未消费

然后执行两个消费者,两个消费者都没有消费到任何消息

并且,控制台页面只是多了2个消费者,已经消费的消息还是0

为什么呢?还记得前面的理论基础说的吗?就是这个原因

继续,我们在两个消费者启动好的前提下,再执行生产者, 这个时候会发现两个消费者都消费了该消息

再看下控制台页面

已消费消息这里是2,这个2并不是说之前发的两个消息都消费了,而是说第二个消息消费了2次, 1 * 2 = 2

不信的话,可以再执行一遍生产者,这个时候就是4,而不是3

累计发送过3条消息,消息消费了4次,这里的4就是后面两条分别被消费了2次, 2 * 2 = 4

三、两种模式比较

好,到这里,发布/订阅模式就介绍完了。

如果有收获,就点个赞呗

原文地址:https://www.cnblogs.com/sam-uncle/p/10990324.html

时间: 2024-10-13 11:48:56

ActiveMQ入门系列三:发布/订阅模式的相关文章

ActiveMQ入门系列二:入门代码实例(点对点模式)

在上一篇<ActiveMQ入门系列一:认识并安装ActiveMQ(Windows下)>中,大致介绍了ActiveMQ和一些概念,并下载.安装.启动他,还访问了他的控制台页面. 这篇,就用代码实例说下如何实现消息的生产和消费. 一.理论基础 同RabbitMQ一样,ActiveMQ中也是有两种模式: 点对点模式(Point to Point,简写为PTP) 发布/订阅模式(Publish & Subscribe,简写为Pub & Sub) 通过上一篇我们知道了制造消息的应用叫生产

ActiveMQ简单简绍(“点对点通讯”和 “发布订阅模式”)

ActiveMQ简单简绍 MQ简介: MQ全称为Message Queue, 消息队列(MQ)是一种应用程序对应用程序的通信方法.应用程序通过写和检索出入列队的针对应用程序的数据(消息)来通信,而无需专用连接来链接它们.消息传递指的是程序之间通过在消息中发送数据进行通信,而不是通过直接调用彼此来通信,直接调用通常是用于诸如远程过程调用的技术.排队指的是应用程序通过队列来通信.队列的使用除去了接收和发送应用程序同时执行的要求.其中较为成熟的MQ产品有IBMWEBSPHERE MQ. MQ特点: M

ActiveMQ发布订阅模式

ActiveMQ的另一种模式就SUB/HUB即发布订阅模式,是SUB/hub就是一拖N的USB分线器的意思.意思就是一个来源分到N个出口.还是上节的例子,当一个订单产生后,后台N个系统需要联动,但有一个前提是都需要收到订单信息,那么我们就需要将一个生产者的消息发布到N个消费者. 生产者: try { //Create the Connection Factory IConnectionFactory factory = new ConnectionFactory("tcp://localhost

python使用rabbitMQ介绍三(发布订阅模式)

一.模式介绍 在前面的例子中,消息直接发送到queue中. 现在介绍的模式,消息发送到exchange中,消费者把队列绑定到exchange上. 发布-订阅模式是把消息广播到每个消费者,每个消费者接收到的消息都是相同的. 一个生产者,多个消费者,每一个消费者都有自己的一个队列,生产者没有将消息直接发送到队列,而是发送到了交换机,每个队列绑定交换机,生产者发送的消息经过交换机,到达队列,实现一个消息被多个消费者获取的目的.需要注意的是,如果将消息发送到一个没有队列绑定的exchange上面,那么该

RabbitMQ学习第三记:发布/订阅模式(Publish/Subscribe)

工作队列模式是直接在生产者与消费者里声明好一个队列,这种情况下消息只会对应同类型的消费者. 举个用户注册的列子:用户在注册完后一般都会发送消息通知用户注册成功(失败).如果在一个系统中,用户注册信息有邮箱.手机号,那么在注册完后会向邮箱和手机号都发送注册完成信息.利用MQ实现业务异步处理,如果是用工作队列的话,就会声明一个注册信息队列.注册完成之后生产者会向队列提交一条注册数据,消费者取出数据同时向邮箱以及手机号发送两条消息.但是实际上邮箱和手机号信息发送实际上是不同的业务逻辑,不应该放在一块处

C# 委托和事件 与 观察者模式(发布-订阅模式)讲解 by天命

使用面向对象的思想 用c#控制台代码模拟猫抓老鼠 我们先来分析一下猫抓老鼠的过程 1.猫叫了 2.所有老鼠听到叫声,知道是哪只猫来了 3.老鼠们逃跑,边逃边喊:"xx猫来了,快跑啊!我是老鼠xxx" 一  双向耦合的代码 首先需要一个猫类Cat 一个老鼠类Rat 和一个测试类Program 老鼠类的代码如下 //老鼠类 public class Rat { public string Name { get; set; } //老鼠的名字 public Cat MyCat { get;

RxJava入门系列三,响应式编程

RxJava入门系列三,响应式编程 在RxJava入门系列一,我向你介绍了RxJava的基础架构.RxJava入门系列二,我向你展示了RxJava提供的多种牛逼操作符.但是你可能仍然没能劝服自己使用RxJava,这一篇博客里我将向你展示RxJava提供的其他优势,没准了解了这些优势,你就真的想去使用RxJava了. 异常处理 直到目前为止,我都没有去介绍onComplete()和onError()方法.这两个方法是用来停止Observable继续发出事件并告知观察者为什么停止(是正常的停止还是因

理解《JavaScript设计模式与开发应用》发布-订阅模式的最终版代码

最近拜读了曾探所著的<JavaScript设计模式与开发应用>一书,在读到发布-订阅模式一章时,作者不仅给出了基本模式的通用版本的发布-订阅模式的代码,最后还做出了扩展,给该模式增加了离线空间功能和命名空间功能,以达到先发布再订阅的功能和防止名称冲突的效果.但是令人感到遗憾的是最终代码并没有给出足够的注释.这让像我一样的小白就感到非常的困惑,于是我将这份最终代码仔细研究了一下,并给出了自己的一些理解,鉴于能力有限,文中观点可能并不完全正确,望看到的大大们不吝赐教,谢谢! 下面是添加了个人注释的

redis的发布订阅模式

概要 redis的每个server实例都维护着一个保存服务器状态的redisServer结构 struct redisServer { /* Pubsub */ // 字典,键为频道,值为链表 // 链表中保存了所有订阅某个频道的客户端 // 新客户端总是被添加到链表的表尾 dict *pubsub_channels;  /* Map channels to list of subscribed clients */ // 这个链表记录了客户端订阅的所有模式的名字 list *pubsub_pa