最近项目中使用到了消息中间件ActiveMQ,现将使用中的一些心得记录下来,不足之处还请各位不吝指出,以便改进。本次入门系列大致分为三部分:JMS相关基本介绍;ActiveMQ基本介绍,安装及入门实例;Spring与ActiveMQ整合实例。下面开始进入第一部分。
什么是Java消息服务(JMS)
根据百度百科和维基百科,JMS即Java消息服务(Java Message Service)应用程序接口,是一个Java平台中关于面向中间件(MOM)的API,用于在两个应用程序之间,或分布式系统中发送消息,进行异步通信。Java消息服务是一个与具体平台无关的API,绝大多数MOM提供商都对JMS提供支持。
The Java Message Service(JMS)API is a Java Message Oriented Middleware(MOM)API for sending messages between two or more clients. It is an implementation to handle the Producer-consumer problem. JMS is a part of the Java Platform, Enterprise Edition, and is defined by a specification developed under the Java Community Process as JSR 914. It is a messagin standard that allows application components based on the Java Enterprise Edition(Java EE)to create, send, receive, and read messages. It allows the communication between different components of a distributed application to be loosely coupled, reliable, and asynchronous.
为什么需要JMS
JMS比较常见的使用场景就是维护两个或多个服务之间的系统通信。比如在某个分布式系统中,有两个服务系统:消费系统A(Oracle)、订单系统B(MySQL),这两个系统可能不在一个服务器甚至不在一个机房,当消费系统A消费了某个订单后,可能希望订单系统B能够对该订单进行某些操作,比如入库等等,这里就可以使用到Java中的消息服务了。在消费系统A消费之后,发送一个包含该订单的消息通知系统B,系统B接收消息然后做相应的业务逻辑处理。还有很多其他的业务场景,都可以使用JMS进行相应的处理。
JMS的优势
JMS使得系统间的通信非常方便,可以进行异步通信,可以对各系统模块之间进行解耦。
JMS的消息传送模型
JMS消息通常有两种类型:点对点或队列消息模型、发布/订阅消息模型。
点对点或队列消息模型
在点对点或队列模型下,一个生产者向一个特定的队列发布消息,一个消费者从该队列中读取消息。这里,生产者知道消费者的队列,并直接将消息发送到消费者的队列。这种模型的特性如下:
(1)、只有一个消费者将获得消息
(2)、生产者不需要在接收者消费该消息期间处于运行状态,接收者也同样不需要在消息发送时处于运行状态
(3)、每一个成功处理的消息都由接收者签收(acknowledgement)

点对点模型
发布/订阅消息模型
发布/订阅模型支持向一个特定的消息主题发布消息。0或多个订阅者可能
对接收来自特定消息主题的消息感兴趣。在这种模型下,发布者和订阅者往往彼此不知道对方。这种模式好比是匿名公告板。这种模型的特性如下:
(1)、一个消息可以传递给多个订阅者
(2)、在发布者和订阅者之间存在时间依赖性。只有当客户端创建订阅后,才能接收消息,且订阅者需要一直保持活动状态以接收消息
(3)、为了可以缓和这种严格的时间依赖性,JMS允许订阅者创建持久化订阅。在这种持久订阅的情况下,在订阅者未连接时发布的消息将在订阅者重新连接时接收到发布者的消息。

发布/订阅模型
接收消息
在JMS中,消息可以使用以下两种方式接收:同步、异步。
同步
消息订阅者可以调用receive()方法来以同步的方式接收消息,在receive()方法中,消息未到达或者在到达指定时间之前,方法会阻塞,直到消息可用。
异步
若需要异步接收,订阅者需要实现MessageListener消息监听器接口,覆盖onMessage()方法,只要消息到达,JMS服务提供者就会调用监听器的onMessage()方法来传递消息。
JMS的编程接口
JMS应用程序由如下基本模块组成:
1、管理对象(Administered objects)-连接工厂(Connection Factories)和目的地(Destination)
2、连接对象(Connections)
3、会话(Sessions)
4、消息生产者(Message Producers)
5、消息消费者(Message Consumers)
6、消息监听者(Message Listeners)

编程接口
下面简要介绍一下各模块:
1、JMS管理对象
管理对象(Administered Object)是预先配置的JMS对象,由系统管理员为使用JMS的客户端创建,主要有两个被管理的对象:
连接工厂(ConnectionFactory)、目的地(Destination)。
这两个管理对象由JMS系统管理员通过使用Application Server管理控制台创建,存储在应用程序服务器的JNDI名字空间或JNDI注册表。
连接工厂(ConnectionFactory):客户端使用一个连接工厂对象连接到JMS服务提供者,它创建了JMS服务提供者和客户端之间的连接。JMS客户端(如消息发送者或接收者)会在JNDI名字空间中搜索并获取该连接。使用该连接,客户端能够与目的地通讯,向队列或者话题发送/接收消息。
目的地(Destination):指消息被发送的目的地以及客户端接收消息的来源。JMS有两种类型的目的地:队列、发布/订阅。
2、JMS连接对象
连接对象封装了与JMS服务提供者之间的虚拟连接。
3、JMS会话对象
Session是一个单线程上下文,用于生产和消费消息,可以创建消息生产者和消息消费者。
4、JMS消息生产者
消息生产者由javax.jms.Session对象创建,用于向目的地发送消息,生产者实现javax.jms.MessageProducer接口,可以为队列、发布/订阅或者目的地创建生产者。
5、JMS消息消费者
同消息生产者类似,也是由javax.jms.Session对象创建,用于接收发向目的地的消息。消费者实现javax.jms.MessageConsumer接口,可以为队列、发布/订阅、或者目的地创建消费者。
6、JMS消息监听器
JMS消息监听器是消息的默认事件处理者,它实现了javax.jms.MessageListener接口,该接口包含一个onMessage(javax.jms.Message message)方法,在该方法中可以处理接收到消息后续的业务逻辑。
代码简例
消息生产者JMSProducer:
import org.apache.activemq.ActiveMQConnectionFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import javax.jms.*;
/** * 消息生产者 * * @author admin * @date 2017/12/16 */
public class JMSProducer {
private static final Logger LOGGER = LoggerFactory.getLogger(JMSProducer.class);
/** * 默认连接用户名 */ private static final String USERNAME = "admin";
/** * 默认连接密码 */ private static final String PASSWORD = "admin";
/** * 默认连接地址 */ private static final String BROKERURL = "tcp://127.0.0.1:61616";
private static final int SENDNUM = 10;
public static void main(String[] args) {
try {
// 连接工厂 ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKERURL);
// 通过连接工厂获取连接 Connection connection = connectionFactory.createConnection();
// 启动连接 connection.start();
// 会话,接收或者发送消息的线程 Session session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
// 消息的目的地:创建一个名称为HelloWorld的消息队列 Destination destination = session.createQueue("HelloWorld");
// 消息生产者 MessageProducer messageProducer = session.createProducer(destination);
// 发送消息 sendMessage(session, messageProducer); } catch (JMSException e) { e.printStackTrace(); LOGGER.error("JMSProducer 发送信息异常!异常信息为:{},e:{}!", e.getMessage(), e); } }
public static void sendMessage(Session session, MessageProducer messageProducer) throws JMSException {
for (int i = 0; i < JMSProducer.SENDNUM; i++) {
// 创建一条文本消息 TextMessage textMessage = session.createTextMessage("ActiveMQ发送消息:" + i); System.out.println("ActiveMQ发送消息: " + i);
// 通过消息生产者发送消息 messageProducer.send(textMessage); } } }
消息消费者JMSConsumer :
import org.apache.activemq.ActiveMQConnectionFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import javax.jms.*;
/** * 消息消费者 * * @author admin * @date 2017/12/16 */
public class JMSConsumer {
private static final Logger LOGGER = LoggerFactory.getLogger(JMSConsumer.class);
/** * 默认连接用户名 */ private static final String USERNAME = "admin";
/** * 默认连接密码 */ private static final String PASSWORD = "admin";
/** * 默认连接地址 */ private static final String BROKERURL = "tcp://127.0.0.1:61616";
public static void main(String[] args) { ConnectionFactory connectionFactory; Connection connection; Session session; Destination destination; MessageConsumer messageConsumer;
// 实例化连接工厂 connectionFactory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKERURL);
try { connection = connectionFactory.createConnection(); session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE); destination = session.createQueue("HelloWorld");
// 创建消息消费者 messageConsumer = session.createConsumer(destination);
while (true) { Message receive = messageConsumer.receive(); TextMessage textMessage = (TextMessage) receive;
if (null != textMessage) { System.out.println("消费者接收到消息: " + textMessage.getText()); } else {
break; } } } catch (JMSException e) { e.printStackTrace(); LOGGER.error("JMSConsumer 接收信息异常!异常信息为:{},e:{}!", e.getMessage(), e); } } }
消息监听器MyMessageListener :
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;
import javax.jms.TextMessage;
/** * 消息监听器 * * @author admin * @date 2018/01/25 */
public class MyMessageListener implements MessageListener {
private static final Logger LOGGER = LoggerFactory.getLogger(MyMessageListener.class);
@Override public void onMessage(Message message) {
try { TextMessage textMessage = (TextMessage) message;
// to do sth... System.out.println("消息监听器监听到的信息是: " + textMessage.getText()); } catch (JMSException e) { e.printStackTrace(); LOGGER.error("MyMessageListener 监听信息异常!异常信息为:{},e:{}!", e.getMessage(), e); } } }
JMS消息结构
JMS客户端使用JMS消息与系统通讯,JMS消息虽然格式简单,但是非常灵活,JMS消息由三部分组成:
1、消息头
JMS消息头预定义了若干字段用于客户端与服务端之间识别和发送信息,预编译头如下:
– JMSDestination : 消息发送的目的地,主要是指Queue和Topic,由session创建而由生产者的send方法设置.
– JMSDeliveryMode:传送模式:有两种即久模式和非持久模式。一条持久性的消息应该被传输"一次仅仅一次",这就意味着如果JMS提供者出现故障,该消息并不会丢失,它会在服务器恢复之后再次传递。一条非持久的消息最多会传递一次,这意味着服务器出现故障,该消息将永远丢失。由session穿件由消息生产者的send方法设置
– JMSMessageID:唯一识别每个消息的标识,由JMS消息生产者产生。由send方法设置
– JMSTimestamp:一个JMS Provider在调用send()方法时自动设置,它是消息被发送和消费者实际接收的时间差。由客户端设置
– JMSCorrelationID:用来连接到另外一个消息,典型的应用是在回复消息中连接到原消息。在大多数情况下,JMSCorrelationID用于将一条消息标记为对JMSMessageID标示的上一条消息的应答,不过,JMSCorrelationID可以是任何值,不仅仅是JMSMessageID。由客户端设置
– JMSReplyTo:提供本消息回复消息的目的地址,由客户端设置
– JMSRedelivered:如果一个客户端收到一个设置了JMSRedelivered属性的消息,则表示可能客户端曾经在早些时候收到过该消息,但并没有签收(acknowledged)。如果该消息被重新传送,JMSRedelivered=true 否则 JMSRedelivered=flase 。由JMS Provider设置
– JMSType:消息类型的标识符,由客户端设置
– JMSExpiration:消息过期时间,等于Destination的send方法中的timeToLive值加上发送时刻的GMT的时间值。如果timeToLive值等于零,则JMSExpiration被设置为零,表示该消息永不过期。如果发送后,在消息过期时间之后消息还没有被发送到目的地,则该消息被清除。由send方法设置
– JMSPriority:消息优先级,从0-9十个级别,0-4是普通消息,5-9是加急消息。JMS不要求JMS Provider严格按照这十个优先级发送消息,但必须保证加急消息要先于普通消息到达,默认是4级。由send方法设置
一个消息的消息头有这些属性,我们可以按照需要对这个消息的消息进行设计,在将这个消息使用消息生产者的send()方法发送到消息服务上。
2、消息属性
我们可以给消息设置自定义属性,这些属性主要是提供给应用程序的。对于实现消息过滤功能,消息属性非常有用,JMS API定义了一些标准属性,JMS服务提供者可以选择性的提供部分标准属性。
3、消息体
在消息体中,JMS API定义了五种类型的消息格式,让我们可以以不同的形式发送和接受消息,并提供了对已有消息格式的兼容。不同的消息类型如下:
Text message : javax.jms.TextMessage,表示一个文本对象。
Object message : javax.jms.ObjectMessage,表示一个JAVA对象。
Bytes message : javax.jms.BytesMessage,表示字节数据。
Stream message :javax.jms.StreamMessage,表示java原始值数据流。
Map message : javax.jms.MapMessage,表示键值对。
以上,是JMS入门的一些理论知识,不足之处,欢迎批评指正。
PS:夜已醉了,夜已醉倒了,让它安静到天晓……欣赏一首《月半弯》,晚安。




