1、 准备工作

1、在我们的linux虚拟机中安装好我们的ActiveMQ,并启动服务,这里我们不详细介绍。

2、关闭防火墙,使得我们能在我们的windows服务器上访问到我们的linux的ActiveMQ,在浏览器中输入127.0.0.1:8161,默认的初始用户名和密码都是admin,访问成功后看到控制台:

在这里插入图片描述在这里插入图片描述
注意:ActiveMQ采用61616端口提供JMS服务,采用8161提供管理控制台服务。

3、创建一个Maven工程,POM.xml导入我们的依赖

<!-- https://mvnrepository.com/artifact/org.apache.activemq/activemq-all -->
        <dependency>
            <groupId>org.apache.activemq</groupId>
            <artifactId>activemq-all</artifactId>
            <version>5.15.11</version>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.xbean/xbean-spring -->
        <dependency>
            <groupId>org.apache.xbean</groupId>
            <artifactId>xbean-spring</artifactId>
            <version>4.15</version>
        </dependency>

2、JMS编码的总体架构

这一点和我们之前的JDBC的编码套路十分类似,我们重点需要知道两种模式的操作。

在这里插入图片描述
JMS开发的基本步骤
1:创建一个connection factory
2:通过connection factory来创建JMS connection
3:启动JMS connection
4:通过JMS connection创建JMS session
5:创建JMS destination(目的地 队列/主题)
6:创建JMS producer或者创建JMS consume并设置destination
7:创建JMS consumer或者注册一个JMS message listener
8:发送(send)或者接收(receive)JMS message
9:关闭所有JMS资源

2、1 queue队列模式

在这里插入图片描述
在点对点的消息传递域中,目的地被称作队列(queue),由图可以知道,我们生产者生产的消息放在队列中,只能被一个消费者消费,并且消费了之后,就不能被其他消费者再次消费。

上手代码:生产消息的生产者

ublic class JmsProduce {
    public static final String ACTIVEMQ_URL="tcp://192.168.126.129:61616";
    public static final String QUEUE_NAME="queue01";
    public static void main(String[] args) throws JMSException {
        //1、获取对应的connectionfactory,按照给定的url地址,采用默认的用户名和密码
        ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL);
        //2、通过连接工厂,获取连接connection,并启动
        Connection connection = activeMQConnectionFactory.createConnection();
        connection.start();
        //3、创建回话session() 有两个参数,第一个是事务,第二个是签收
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        //4、创建目的地(是队列还是主题topic)
//        Destination destination = session.createQueue(QUEUE_NAME);
        Queue queue = session.createQueue(QUEUE_NAME);
        //5、创建消息的生产者
        MessageProducer messageProducer = session.createProducer(queue);
        //6、通过使用messageProducer生产3条消息到MQ的队列里面

        for (int i = 1; i < 7; i++) {
            //7、创建消息 文本消息
            TextMessage textMessage = session.createTextMessage("msg---:" + i);
            //8、通过MessageProducer发送给MQ
            messageProducer.send(textMessage);
        }
        //9、关闭资源
        messageProducer.close();
        connection.close();
        System.out.println("**********消息发送到mq成功**********");

    }
}

消费消息的消费者:

public class JmsConsumer1 {
    public static final String ACTIVEMQ_URL = "tcp://192.168.126.129:61616";
    public static final String QUEUE_NAME = "queue01";

    public static void main(String[] args) throws JMSException, IOException {
        System.out.println("我是1号消费者");
        //1、获取对应的connectionfactory,按照给定的url地址,采用默认的用户名和密码
        ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL);
        //2、通过连接工厂,获取连接connection,并启动
        Connection connection = activeMQConnectionFactory.createConnection();
        connection.start();
        //3、创建回话session() 有两个参数,第一个是事务,第二个是签收
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        //4、创建目的地(是队列还是主题topic)
        // Destination destination = session.createQueue(QUEUE_NAME);
        Queue queue = session.createQueue(QUEUE_NAME);
        //5、创建消费者
        MessageConsumer consumer = session.createConsumer(queue);
       /* 同步阻塞方式(receive())
       订阅者或接受者调用MessageConsumer的receive()方法来接受消息,这个方法可以传递一个参数,设置过期时间,在这之前将一致处于阻塞状态
        while (true) {
            TextMessage message = (TextMessage) consumer.receive();
            if(message!=null){
                System.out.println("******消费者接收到消息:"+message.getText());
            }else break;
        }
        consumer.close();
        connection.close();*/
        consumer.setMessageListener(new MessageListener(){

            public void onMessage(Message message) {
                if(message!=null && message instanceof TextMessage){
                    TextMessage textMessage = (TextMessage) message;
                    try {
                        System.out.println("消费者接收到消息"+textMessage.getText());
                    } catch (JMSException e) {
                        e.printStackTrace();
                    }
                }
            }
        });
        System.in.read();
        consumer.close();
        session.close();
        connection.close();
    }
}

控制台:
在这里插入图片描述
Number Of Pending Messages=等待消费的消息,这个是未出队列的数量,公式=总接收数-总出队列数。
Number Of Consumers=消费者数量,消费者端的消费者数量。
Messages Enqueued=进队消息数,进队列的总消息量,包括出队列的。这个数只增不减。
Messages Dequeued=出队消息数,可以理解为是消费者消费掉的数量。
总结:
当有一个消息进入这个队列时,等待消费的消息是1,进入队列的消息是1。
当消息消费后,等待消费的消息是0,进入队列的消息是1,出队列的消息是1。
当再来一条消息时,等待消费的消息是1,进入队列的消息就是2。

2、2 topic主题

在这里插入图片描述
topic发布的信息被多个消费者同时消费但是如果提前没有消费者,那么生产的消息就是一个废消息,不会被任何消费者消费。

生产者代码:

和之前的queue的生产者的基本流程一样。

public class JmsProduce {
    public static final String ACTIVEMQ_URL = "tcp://192.168.126.129:61616";
    public static final String QUEUE_NAME = "topic_atLizy";

    public static void main(String[] args) throws JMSException {
        ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL);
        Connection connection = connectionFactory.createConnection();
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Topic topic = session.createTopic(QUEUE_NAME);
        MessageProducer producer = session.createProducer(topic);
        for (int i = 0; i < 3; i++) {
            TextMessage textMessage = session.createTextMessage("topic发送消息:top——" + i);
            producer.send(textMessage);
        }
        producer.close();
        session.close();
        connection.close();
    }
}

消费者代码
和queue的消费者也基本类似:

public class JmsConsumer1 {
    public static final String ACTIVEMQ_URL = "tcp://192.168.126.129:61616";
    public static final String QUEUE_NAME = "topic_atLizy";

    public static void main(String[] args) throws JMSException, IOException {
        System.out.println("我是3号消费者");
        //1、获取对应的connectionfactory,按照给定的url地址,采用默认的用户名和密码
        ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(ACTIVEMQ_URL);
        //2、通过连接工厂,获取连接connection,并启动
        Connection connection = activeMQConnectionFactory.createConnection();
        connection.start();
        //3、创建回话session() 有两个参数,第一个是事务,第二个是签收
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        //4、创建目的地(是队列还是主题topic)
        // Destination destination = session.createQueue(QUEUE_NAME);
        Topic topic = session.createTopic(QUEUE_NAME);
        //5、创建消费者
        MessageConsumer consumer = session.createConsumer(topic);
       /* 同步阻塞方式(receive())
       订阅者或接受者调用MessageConsumer的receive()方法来接受消息,这个方法可以传递一个参数,设置过期时间,在这之前将一致处于阻塞状态
        while (true) {
            TextMessage message = (TextMessage) consumer.receive();
            if(message!=null){
                System.out.println("******消费者接收到消息:"+message.getText());
            }else break;
        }
        consumer.close();
        connection.close();*/
        consumer.setMessageListener(new MessageListener(){

            public void onMessage(Message message) {
                if(message!=null && message instanceof TextMessage){
                    TextMessage textMessage = (TextMessage) message;
                    try {
                        System.out.println("消费者接收到消息"+textMessage.getText());
                    } catch (JMSException e) {
                        e.printStackTrace();
                    }
                }
            }
        });
        System.in.read();
        consumer.close();
        session.close();
        connection.close();
    }
}

控制台:
在这里插入图片描述

2、3 总结

在这里插入图片描述

更多推荐