消息队列入门(四)ActiveMQ的应用实例

简介:

一、部署和启动ActiveMQ

去官网下载:http://activemq.apache.org/

我下载的是apache-activemq-5.12.0-bin.tar.gz,

解压到本地目录,进入到bin路径下,
运行activemq启动ActiveMQ。

运行方式:
启动 ./activemq start

ActiveMQ默认使用的TCP连接端口是61616,
5.0以上版本默认启动时,开启了内置的Jetty服务器,可以进入控制台查看管理。

启动ActiveMQ以后,登陆:http://localhost:8161/admin/

默认用户名admin/admin

这里我在虚拟机里启动,访问地址:
http://192.168.106.128:8161/admin/

ActiveMQ的控制台功能十分强大,管理起来也很直观。

二、使用Java连接

1.创建POM文件

在Eclipse中新建Java工程,这里使用Maven管理依赖,
下面是pom.xml:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
< project  xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
< modelVersion >4.0.0</ modelVersion >
< groupId >activemq-sample</ groupId >
< artifactId >activemq-sample</ artifactId >
< version >0.0.1-SNAPSHOT</ version >
< name >activemq-sample</ name >
< description >an activemq practice</ description >
< build >
< sourceDirectory >src</ sourceDirectory >
< plugins >
 
< plugin >
< artifactId >maven-compiler-plugin</ artifactId >
< version >3.1</ version >
< configuration >
< source >1.7</ source >
< target >1.7</ target >
</ configuration >
</ plugin >
<!-- activemq-core 5.7.0 使用bunble打包,需要添加相关插件 -->
< plugin >
< groupId >org.apache.felix</ groupId >
< artifactId >maven-bundle-plugin</ artifactId >
< extensions >true</ extensions >
</ plugin >
 
</ plugins >
</ build >
< dependencies >
<!-- activemq的maven依赖 -->
< dependency >
< groupId >org.apache.activemq</ groupId >
< artifactId >activemq-core</ artifactId >
< version >5.7.0</ version >
< type >bundle</ type >
</ dependency >
 
</ dependencies >
 
</ project >

  

在第一次添加activemq的maven依赖时报错,后来发现activemq-core 5.7.0采用了bundle的打包方式,

必须在pom中配置maven-bundle-plugin。

2.创建消息创建者 MsgProducer:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
import  javax.jms.Connection;
import  javax.jms.ConnectionFactory;
import  javax.jms.Destination;
import  javax.jms.JMSException;
import  javax.jms.MessageProducer;
import  javax.jms.Session;
import  javax.jms.TextMessage;
 
import  org.apache.activemq.ActiveMQConnectionFactory;
 
/**
* @Description: Message Producer
* @author: Bing Yue
*/
public  class  MsgProducer {
//如果你在本地启动,可以直接使用空的ActiveMQConnectionFactory构造函数
private  static  final  String BROKER_URL= "failover://tcp://192.168.106.128:61616" ;
 
public  static  void  main(String[] args)  throws  JMSException, InterruptedException{
//创建连接工厂
ConnectionFactory connectionFactory= new  ActiveMQConnectionFactory(BROKER_URL);
//获得连接
Connection conn = connectionFactory.createConnection();
//start
conn.start();
 
//创建Session,此方法第一个参数表示会话是否在事务中执行,第二个参数设定会话的应答模式
Session session = conn.createSession( false , Session.AUTO_ACKNOWLEDGE);
 
//创建队列
Destination dest = session.createQueue( "test-queue" );
//创建消息生产者
MessageProducer producer = session.createProducer(dest);
 
for  ( int  i= 0 ;i< 100 ;i++) {
//初始化一个mq消息
TextMessage message = session.createTextMessage( "这是第 "  + i+ " 条消息!" );
//发送消息
producer.send(message);
System.out.println( "send message:消息" +i);
//暂停3秒
Thread.sleep( 3000 );
}
 
//关闭mq连接
conn.close();
}
 
}

  

3.创建消息接收者 MsgProducer:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
import  javax.jms.Connection;
import  javax.jms.ConnectionFactory;
import  javax.jms.Destination;
import  javax.jms.JMSException;
import  javax.jms.Message;
import  javax.jms.MessageConsumer;
import  javax.jms.MessageListener;
import  javax.jms.Session;
import  javax.jms.TextMessage;
 
import  org.apache.activemq.ActiveMQConnectionFactory;
 
/**
*
* @Description: Message Consumer
* @author: Bing Yue
*/
public  class  MsgConsumer  implements  MessageListener {
 
private  static  final  String BROKER_URL= "failover://tcp://192.168.106.128:61616" ;
 
public  static  void  main(String[] args)  throws  JMSException{
 
//创建连接工厂
ConnectionFactory connectionFactory= new  ActiveMQConnectionFactory(BROKER_URL);
//获得连接
Connection conn = connectionFactory.createConnection();
//start
conn.start();
 
//创建Session,此方法第一个参数表示会话是否在事务中执行,第二个参数设定会话的应答模式
Session session = conn.createSession( false , Session.AUTO_ACKNOWLEDGE);
//创建队列
Destination dest = session.createQueue( "test-queue" );
//创建消息生产者
MessageConsumer consumer = session.createConsumer(dest);
 
//初始化MessageListener
MsgConsumer msgConsumer =  new  MsgConsumer();
 
//给消费者设定监听对象
consumer.setMessageListener(msgConsumer);
 
}
 
 
/**
* 消费者需要实现MessageListener接口
* 接口有一个onMessage(Message message)需要在此方法中做消息的处理
*/
@Override
public  void  onMessage(Message msg) {
TextMessage txtMessage = (TextMessage)msg;
try  {
System.out.println( "get message:"  + txtMessage.getText());
catch  (JMSException e) {
e.printStackTrace();
}
 
}
}

  

运行MsgProducer,

登录后台查看test-queue队列,可以看到发出的消息正在等待被处理:

 

运行MsgConsumer,接收消息并在控制台打印:

通过这个实例可以对ActiveMQ的应用有一个简单的了解。

代码地址:https://github.com/bingyue/activemq-sample

在实际开发中,通常还需要设置优先级处理,大部分情况下,消息的发送和接收方都会启用多线程,
通过线程池来提高处理效率,解耦的同时保持业务处理能力。

 


本文转自邴越博客园博客,原文链接:http://www.cnblogs.com/binyue/p/4763767.html,如需转载请自行联系原作者

相关文章
|
3月前
|
消息中间件 NoSQL Java
Redis Streams在Spring Boot中的应用:构建可靠的消息队列解决方案【redis实战 二】
Redis Streams在Spring Boot中的应用:构建可靠的消息队列解决方案【redis实战 二】
246 1
|
2月前
|
消息中间件 Linux API
Linux进程间通信(IPC) Linux消息队列:讲解POSIX消息队列在Linux系统进程间通信中的应用和实践
Linux进程间通信(IPC) Linux消息队列:讲解POSIX消息队列在Linux系统进程间通信中的应用和实践
27 1
Linux进程间通信(IPC) Linux消息队列:讲解POSIX消息队列在Linux系统进程间通信中的应用和实践
|
3月前
|
消息中间件 存储 负载均衡
简单入门:消息队列的概念和应用
在复杂的系统架构中,组件间的通信是至关重要的问题。消息队列作为一种解决方案,能够使组件之间的通信更加高效、可靠。本文将从简单到复杂,逐步向您介绍消息队列的概念、使用场景以及如何实现。
100 3
|
4月前
|
消息中间件 存储 Kafka
MQ消息队列学习入门
MQ消息队列学习入门
77 0
|
4月前
|
消息中间件 监控 负载均衡
Kafka高级应用:如何配置处理MQ百万级消息队列?
在大数据时代,Apache Kafka作为一款高性能的分布式消息队列系统,广泛应用于处理大规模数据流。本文将深入探讨在Kafka环境中处理百万级消息队列的高级应用技巧。
178 0
|
5月前
|
消息中间件 网络协议 Java
RabbitMQ消息队列基础详解与安装实例
RabbitMQ消息队列基础详解与安装实例
130 0
|
6月前
|
消息中间件 Go 流计算
Golang微服务框架Kratos应用NATS消息队列详解
Golang微服务框架Kratos应用NATS消息队列详解
|
6月前
|
消息中间件 Kafka Go
Golang微服务框架Kratos应用Kafka消息队列
Apache Kafka 是一个分布式数据流处理平台,可以实时发布、订阅、存储和处理数据流。它旨在处理多种来源的数据流,并将它们交付给多个消费者。简而言之,它可以移动大量数据,不仅是从 A 点移到 B 点,而是能从 A 到 Z 的多个点移到任何您想要的位置,并且可以同时进行。
122 0
|
6月前
|
消息中间件 网络协议 物联网
Golang微服务框架Kratos应用MQTT消息队列
MQTT 协议 是由`IBM`的`Andy Stanford-Clark博士`和`Arcom`(已更名为Eurotech)的`Arlen Nipper博士`于 1999 年发明,用于石油和天然气行业。工程师需要一种协议来实现最小带宽和最小电池损耗,以通过卫星监控石油管道。最初,该协议被称为消息队列遥测传输,得名于首先支持其初始阶段的 IBM 产品 MQ 系列。2010 年,IBM 发布了 MQTT 3.1 作为任何人都可以实施的免费开放协议,然后于 2013 年将其提交给结构化信息标准促进组织 (OASIS) 规范机构进行维护。2019 年,OASIS 发布了升级的 MQTT 版本 5。
45 0
|
6月前
|
消息中间件 Go 网络性能优化
Golang微服务框架Kratos应用NATS消息队列
NATS是由CloudFoundry的架构师Derek开发的一个开源的、轻量级、高性能的,支持发布、订阅机制的分布式消息队列系统。它的核心基于EventMachine开发,代码量不多,可以下载下来慢慢研究。其核心原理就是基于消息发布订阅机制。每个台服务 器上的每个模块会根据自己的消息类别,向MessageBus发布多个消息主题;而同时也向自己需要交互的模块,按照需要的信息内容的消息主题订阅消息。 NATS原来是使用Ruby编写,可以实现每秒150k消息,后来使用Go语言重写,能够达到每秒8-11百万个消息,整个程序很小只有3M Docker image
89 0

热门文章

最新文章