ActiveMQ的helloword(简单使用)
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 总结

更多推荐


所有评论(0)