背景
本文是《JavaEE 后端从小白到大神》修仙系列第六篇,正式进入JavaEE后端世界。若想详细学习请点击首篇博文,我们开始吧。
第六篇:JMS(消息服务)与异步通信机制
- JMS API 规范;
- 点对点(Queue)与发布订阅(Topic)模型;
- Message、Producer、Consumer;
- Message Listener;
- 事务性消息;
- 与现代 MQ(Kafka、RabbitMQ、ActiveMQ)的对比。
一、JMS 概述
1. 什么是 JMS
JMS(Java Message Service,Java消息服务)是 Jakarta EE(原 Java EE)中定义的消息中间件标准 API,用于在 Java 应用程序之间进行异步、可靠的消息通信。
一句话总结:JMS = Java 平台消息通信的标准 API,让不同的 Java 应用可以通过消息进行异步通信。
2. 为什么需要消息服务?
在分布式系统中,应用之间需要通信,但直接调用(如 RPC、HTTP)存在以下问题:
- 同步阻塞:调用方必须等待被调用方返回结果,期间线程被占用;
- 紧耦合:调用方必须知道被调用方的地址和接口;
- 单点故障:被调用方不可用时,调用方直接失败;
- 流量冲击:突发高并发时,被调用方可能被压垮。
消息服务通过异步、解耦、削峰的方式解决这些问题:
| 问题 |
消息服务的解决方案 |
| 同步阻塞 |
发送方发完消息即可继续,无需等待处理结果 |
| 紧耦合 |
发送方只发消息到目的地,不关心谁消费 |
| 单点故障 |
消息持久化,消费者恢复后可继续消费 |
| 流量冲击 |
消息队列缓冲请求,消费者按自身能力处理 |
3. JMS 的演进历程
| 版本 |
发布时间 |
核心变化 |
| JMS 1.0 |
1998 |
初始规范,定义了两种独立的 API(Queue 和 Topic 各有一套) |
| JMS 1.1 |
2002 |
引入统一 API,一套接口同时支持 Queue 和 Topic |
| JMS 2.0 |
2013(JSR 343) |
引入简化 API(Simplified API),接口更少、使用更方便 |
Jakarta EE 中的对应规范称为 Jakarta Messaging。
4. JMS 核心架构
1
2
3
4
5
6
7
|
┌─────────────┐ ┌─────────────────────────────────────┐ ┌─────────────┐
│ 生产者 │ │ JMS 提供者(Broker) │ │ 消费者 │
│ (Producer) │────▶│ ┌─────────┐ ┌───────────────┐ │────▶│ (Consumer) │
│ │ │ │ Queue │ │ Topic │ │ │ │
│ │ │ │(点对点) │ │ (发布/订阅) │ │ │ │
└─────────────┘ │ └─────────┘ └───────────────┘ │ └─────────────┘
└─────────────────────────────────────┘
|
JMS 架构中的核心角色:
| 角色 |
说明 |
| JMS 提供者(Provider) |
实现 JMS 规范的消息中间件(如 ActiveMQ、Artemis) |
| JMS 客户端 |
使用 JMS API 发送或接收消息的 Java 应用 |
| 消息(Message) |
在客户端之间传递的数据载体 |
| 目的地(Destination) |
消息发送的目标(Queue 或 Topic) |
二、两种消息模型
JMS 定义了两种消息模型:点对点(Point-to-Point) 和 发布/订阅(Publish/Subscribe) 。
1. 点对点模型(PTP——Queue)
1.1 核心特征
- 一个生产者 → 一个队列 → 一个消费者:每条消息只能被一个消费者消费;
- 消息被消费后即从队列移除:不会被其他消费者收到;
- 消费者无需在线:生产者发送消息时,消费者可以不在线,消息会保存在队列中,等消费者上线后获取。
1.2 类比理解
点对点模型就像快递柜:快递员(生产者)把包裹(消息)放进柜子(队列),收件人(消费者)有空时来取。包裹被取走后,柜子里就没有了。快递员不需要等收件人在场。
1.3 适用场景
- 任务分发(一个任务只需要一个 worker 处理);
- 订单处理(一个订单只需要一个系统处理);
- 日志收集(每条日志只需被一个消费者处理一次)。
2. 发布/订阅模型(Pub/Sub——Topic)
2.1 核心特征
- 一个生产者 → 一个主题 → 多个消费者:一条消息可以被多个订阅者同时消费;
- 消息会被所有订阅者收到:每个订阅者都能收到完整消息;
- 消费者需要在线(除非是持久订阅):消息发布时,不在线的非持久订阅者收不到消息。
2.2 类比理解
发布/订阅模型就像微信公众号:公众号(生产者)发一篇文章(消息)到平台(主题),所有关注了该公众号的粉丝(订阅者)都能看到。每个粉丝独立收到完整内容。
2.3 适用场景
- 广播通知(系统公告、配置更新);
- 事件分发(用户注册后,多个系统需要响应);
- 实时数据推送(股票行情、新闻推送)。
3. 两种模型对比
| 对比维度 |
点对点(Queue) |
发布/订阅(Topic) |
| 消息去向 |
一个消费者 |
多个消费者(所有订阅者) |
| 消息是否可重复消费 |
否(消费即移除) |
是(每个订阅者都收到一份) |
| 消费者是否需要在线 |
否(消息持久化在队列中) |
持久订阅需要在线,非持久订阅不需要 |
| 典型场景 |
任务队列、订单处理 |
广播通知、事件驱动 |
核心区别一句话:Queue 是一对一(一条消息只给一个人),Topic 是一对多(一条消息给所有人)。
三、JMS 核心 API
1. JMS 1.1 统一 API(传统 API)
JMS 1.1 引入了一套统一的 API,同时支持 Queue 和 Topic:
| 接口 |
作用 |
ConnectionFactory |
创建 Connection 的连接工厂 |
Connection |
代表与 JMS 提供者的物理连接 |
Session |
发送/接收消息的上下文(单线程) |
Destination |
消息目的地(Queue 或 Topic) |
MessageProducer |
消息生产者 |
MessageConsumer |
消息消费者 |
Message |
消息本身 |
2. JMS 2.0 简化 API(Simplified API)
JMS 2.0 引入的简化 API 减少了接口数量,使用更方便:
| 接口 |
作用 |
ConnectionFactory |
连接工厂(不变) |
JMSContext |
合并了 Connection + Session,是简化的核心入口 |
JMSConsumer |
合并了 MessageConsumer + 部分 Session 功能 |
JMSProducer |
合并了 MessageProducer + 部分 Session 功能 |
JMS 2.0 简化 API 的优势:接口更少、代码更简洁、资源自动管理。
3. 消息类型
JMS 定义了五种消息类型,覆盖了主流的数据格式:
| 消息类型 |
说明 |
适用场景 |
TextMessage |
字符串消息(最常用) |
JSON、XML、普通文本 |
MapMessage |
键值对集合 |
传输结构化的名值对数据 |
BytesMessage |
字节数组 |
二进制数据(文件、图片) |
StreamMessage |
Java 基本类型流 |
顺序传输基本类型数据 |
ObjectMessage |
可序列化的 Java 对象 |
传输复杂对象(需实现 Serializable) |
最佳实践:大多数场景使用 TextMessage 传递 JSON 字符串,既灵活又跨语言兼容。
4. 消息头与属性
每条 JMS 消息都包含消息头(Header) 和可选的属性(Property):
标准消息头(由 JMS 自动设置):
| 消息头 |
说明 |
JMSDestination |
消息的目的地(Queue 或 Topic) |
JMSDeliveryMode |
持久化模式(PERSISTENT / NON_PERSISTENT) |
JMSExpiration |
过期时间 |
JMSPriority |
优先级(0-9,9 最高) |
JMSMessageID |
消息唯一 ID |
JMSTimestamp |
消息发送时间戳 |
JMSCorrelationID |
关联 ID(用于请求-响应模式) |
JMSReplyTo |
回复目的地 |
JMSType |
消息类型标识 |
JMSRedelivered |
是否被重新投递 |
自定义属性:开发者可以添加任意键值对属性,用于过滤和路由。
四、同步消费与异步消费
1. 同步消费(receive())
消费者主动调用 receive() 方法,阻塞等待消息到达。
1
2
3
4
5
6
7
|
// 同步接收(阻塞)
Message message = consumer.receive(); // 一直等待
// 或设置超时时间
Message message = consumer.receive(5000); // 最多等待5秒
if (message != null) {
// 处理消息
}
|
2. 异步消费(MessageListener)
消费者注册一个 MessageListener,消息到达时自动回调 onMessage() 方法。
1
2
3
4
5
6
7
|
consumer.setMessageListener(new MessageListener() {
@Override
public void onMessage(Message message) {
// 消息到达时自动调用
// 处理消息...
}
});
|
异步消费的优势是消费者线程不会被阻塞,可以持续处理其他任务。
3. 两种消费方式对比
| 对比维度 |
同步消费(receive) |
异步消费(MessageListener) |
| 线程模型 |
调用线程阻塞等待 |
消息到达时回调,线程不阻塞 |
| 实时性 |
取决于轮询频率 |
消息到达即处理 |
| 复杂度 |
简单 |
稍复杂(需处理并发) |
| 适用场景 |
简单测试、批量处理 |
生产环境、实时处理 |
五、事务性消息
JMS 支持两种事务模式:本地事务和分布式事务(JTA/XA) 。
1. 本地事务(Local Transaction)
在单个 JMS Session 范围内的事务控制。通过 Session 的 commit() 和 rollback() 管理。
1
2
3
4
5
6
7
8
9
10
11
12
|
Session session = connection.createSession(true, Session.SESSION_TRANSACTED);
// 创建生产者、发送消息...
try {
// 发送多条消息
producer.send(message1);
producer.send(message2);
// 提交事务——两条消息同时成功
session.commit();
} catch (Exception e) {
// 回滚事务——两条消息都不发送
session.rollback();
}
|
关键特性:
- 事务提交前,消息不会真正发送到目的地;
- 事务回滚时,所有消息都不会发送;
- 仅涉及 JMS 单一资源。
2. 分布式事务(JTA/XA)
当一次操作涉及多个资源(如 JMS + 数据库)时,需要分布式事务保证一致性。
典型场景:订单服务收到请求后,需要同时:
- 向数据库写入订单记录;
- 发送 JMS 消息通知库存系统扣减库存。
如果数据库写入成功但消息发送失败,会导致数据不一致。分布式事务通过 两阶段提交(2PC) 保证要么都成功,要么都回滚。
1
2
3
4
5
6
7
8
9
10
11
12
|
// 使用 JTA 分布式事务(伪代码示意)
@Transactional // 容器管理分布式事务
public void createOrder(Order order) {
// 1. 写入数据库(参与分布式事务)
entityManager.persist(order);
// 2. 发送 JMS 消息(参与同一个分布式事务)
jmsProducer.send(orderQueue, orderMessage);
// 方法正常结束 → 两阶段提交,数据库和消息同时提交
// 方法抛出异常 → 两阶段回滚,数据库和消息同时回滚
}
|
JMS 通过实现 JTA 的 XAResource 接口来参与分布式事务。
3. 事务与消息确认(Acknowledgment)
在非事务模式下,消息的确认(Acknowledgment)机制决定了消费者何时通知 Broker 消息已被处理:
| 确认模式 |
说明 |
AUTO_ACKNOWLEDGE |
消息被接收后自动确认(默认) |
CLIENT_ACKNOWLEDGE |
手动调用 message.acknowledge() 确认 |
DUPS_OK_ACKNOWLEDGE |
延迟确认,允许重复消息(性能更高) |
六、消息可靠性保证
1. 持久化(Persistence)
消息可以设置为持久化(PERSISTENT)或非持久化(NON_PERSISTENT):
- 持久化:消息被写入磁盘,Broker 重启后消息不丢失;
- 非持久化:消息仅在内存中,Broker 重启后消息丢失。
1
2
3
4
|
// 设置消息持久化
producer.setDeliveryMode(DeliveryMode.PERSISTENT);
// 或单独设置每条消息
message.setJMSDeliveryMode(DeliveryMode.PERSISTENT);
|
2. 消息过期(Expiration)
消息可以设置过期时间,过期后 Broker 自动丢弃:
1
2
|
// 设置消息 60 秒后过期
producer.setTimeToLive(60000);
|
3. 消息优先级(Priority)
优先级 0-9,9 最高,Broker 优先投递高优先级消息:
1
|
producer.setPriority(9);
|
4. 可靠投递的层次
| 可靠性级别 |
配置 |
说明 |
| 最多一次(At Most Once) |
非持久化 + AUTO_ACK |
消息可能丢失,性能最高 |
| 至少一次(At Least Once) |
持久化 + CLIENT_ACK |
消息不丢,但可能重复 |
| 恰好一次(Exactly Once) |
持久化 + 分布式事务 |
消息不丢不重,性能最低 |
七、实战案例:订单异步处理系统
1. 业务需求
构建一个订单处理系统,包含以下流程:
- 用户下单 → 订单服务接收请求,保存订单到数据库;
- 订单服务发送 JMS 消息通知库存服务扣减库存;
- 库存服务异步消费消息,扣减库存;
- 扣减完成后,发送消息通知订单服务更新订单状态;
- 订单服务收到状态更新消息后,更新订单状态。
2. 架构图
1
2
3
4
5
6
7
8
9
10
11
12
13
|
┌────────────┐ ┌──────────────────────────────────────────────────────┐ ┌────────────┐
│ 用户 │ │ 订单系统 │ │ 库存系统 │
│ 下单 │────▶│ ┌─────────────┐ ┌──────────────────────────┐ │ │ │
└────────────┘ │ │ 订单服务 │───▶│ Queue: order.created │ │────▶│ 扣减库存 │
│ │ (保存订单) │ └──────────────────────────┘ │ │ │
│ └─────────────┘ │ └──────┬─────┘
│ ▲ │ │
│ │ │ │
│ ┌──────┴─────┐ ┌──────────────────────────┐ │ │
│ │ 订单服务 │◀───│ Queue: inventory.updated │ │◀───────────┘
│ │(更新状态) │ └──────────────────────────┘ │
│ └────────────┘ │
└──────────────────────────────────────────────────────┘
|
3. 依赖配置
1
2
3
4
5
6
7
8
9
10
11
12
13
14
|
<!-- Jakarta EE 9+ 使用 jakarta.jms -->
<dependency>
<groupId>jakarta.jms</groupId>
<artifactId>jakarta.jms-api</artifactId>
<version>3.1.0</version>
<scope>provided</scope>
</dependency>
<!-- ActiveMQ Artemis 客户端(JMS 实现) -->
<dependency>
<groupId>org.apache.activemq</groupId>
<artifactId>artemis-jms-client</artifactId>
<version>2.31.2</version>
</dependency>
|
4. 订单服务——发送订单创建消息
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
|
import jakarta.ejb.Stateless;
import jakarta.ejb.TransactionAttribute;
import jakarta.ejb.TransactionAttributeType;
import jakarta.jms.*;
import jakarta.annotation.Resource;
import jakarta.persistence.EntityManager;
import jakarta.persistence.PersistenceContext;
@Stateless
public class OrderService {
@PersistenceContext
private EntityManager em;
@Resource(lookup = "java:jboss/queue/orderQueue")
private Queue orderQueue;
@Resource(lookup = "java:jboss/ConnectionFactory")
private ConnectionFactory connectionFactory;
@TransactionAttribute(TransactionAttributeType.REQUIRED)
public void createOrder(Order order) throws Exception {
// 1. 保存订单到数据库
em.persist(order);
// 2. 发送 JMS 消息通知库存系统
try (JMSContext context = connectionFactory.createContext(
JMSContext.AUTO_ACKNOWLEDGE)) {
JMSProducer producer = context.createProducer();
// 设置持久化
producer.setDeliveryMode(DeliveryMode.PERSISTENT);
// 创建消息(使用 TextMessage 传递 JSON)
TextMessage message = context.createTextMessage();
String orderJson = String.format(
"{\"orderId\":%d,\"productId\":%d,\"quantity\":%d}",
order.getId(), order.getProductId(), order.getQuantity()
);
message.setText(orderJson);
// 发送消息
producer.send(orderQueue, message);
System.out.println("订单消息已发送,订单ID:" + order.getId());
}
}
}
|
5. 库存服务——消费订单消息(MDB)
使用消息驱动 Bean(MDB)异步消费消息,配合上一章的 EJB 知识:
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 jakarta.ejb.MessageDriven;
import jakarta.jms.Message;
import jakarta.jms.MessageListener;
import jakarta.jms.TextMessage;
import jakarta.jms.JMSException;
import jakarta.ejb.ActivationConfigProperty;
import jakarta.persistence.EntityManager;
import jakarta.persistence.PersistenceContext;
import jakarta.ejb.EJB;
@MessageDriven(
activationConfig = {
@ActivationConfigProperty(
propertyName = "destinationType",
propertyValue = "jakarta.jms.Queue"
),
@ActivationConfigProperty(
propertyName = "destination",
propertyValue = "queue/orderQueue"
),
@ActivationConfigProperty(
propertyName = "acknowledgeMode",
propertyValue = "Auto-acknowledge"
)
}
)
public class InventoryConsumerMDB implements MessageListener {
@PersistenceContext
private EntityManager em;
@EJB
private InventoryService inventoryService;
@EJB
private OrderStatusNotifier notifier; // 发送状态更新消息
@Override
public void onMessage(Message message) {
try {
if (message instanceof TextMessage) {
TextMessage textMsg = (TextMessage) message;
String orderJson = textMsg.getText();
System.out.println("库存服务收到订单消息:" + orderJson);
// 解析 JSON,扣减库存
// 假设解析得到 orderId, productId, quantity
// inventoryService.deductStock(productId, quantity);
// 扣减成功后,发送状态更新消息
notifier.sendOrderUpdated(orderId, "INVENTORY_DEDUCTED");
System.out.println("库存扣减完成,订单ID:" + orderId);
}
} catch (JMSException e) {
System.err.println("处理订单消息失败:" + e.getMessage());
// 抛出 RuntimeException 触发消息重发
throw new RuntimeException("库存扣减失败", e);
}
}
}
|
6. 订单状态更新服务(发送状态消息)
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
|
import jakarta.ejb.Stateless;
import jakarta.jms.*;
import jakarta.annotation.Resource;
@Stateless
public class OrderStatusNotifier {
@Resource(lookup = "java:jboss/queue/orderStatusQueue")
private Queue statusQueue;
@Resource(lookup = "java:jboss/ConnectionFactory")
private ConnectionFactory connectionFactory;
public void sendOrderUpdated(Long orderId, String status) {
try (JMSContext context = connectionFactory.createContext(
JMSContext.AUTO_ACKNOWLEDGE)) {
TextMessage message = context.createTextMessage();
message.setText(String.format(
"{\"orderId\":%d,\"status\":\"%s\"}",
orderId, status
));
context.createProducer().send(statusQueue, message);
System.out.println("订单状态更新消息已发送:" + orderId + " -> " + status);
}
}
}
|
7. 订单服务——消费状态更新消息
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
|
@MessageDriven(
activationConfig = {
@ActivationConfigProperty(
propertyName = "destinationType",
propertyValue = "jakarta.jms.Queue"
),
@ActivationConfigProperty(
propertyName = "destination",
propertyValue = "queue/orderStatusQueue"
)
}
)
public class OrderStatusConsumerMDB implements MessageListener {
@PersistenceContext
private EntityManager em;
@Override
public void onMessage(Message message) {
try {
if (message instanceof TextMessage) {
TextMessage textMsg = (TextMessage) message;
String json = textMsg.getText();
// 解析 JSON,更新订单状态
// order.setStatus(status);
// em.merge(order);
System.out.println("订单状态已更新:" + json);
}
} catch (JMSException e) {
throw new RuntimeException("更新订单状态失败", e);
}
}
}
|
八、JMS 与现代消息中间件的对比
1. 主流消息中间件概览
| 消息中间件 |
协议 |
特点 |
适用场景 |
| ActiveMQ Classic |
JMS、OpenWire |
JMS 规范实现,支持事务、XA |
传统 Java 企业应用 |
| ActiveMQ Artemis |
JMS、AMQP、MQTT |
ActiveMQ 的下一代,性能更好 |
Jakarta EE 集成、高性能消息 |
| RabbitMQ |
AMQP |
Erlang 开发,灵活路由,高可用 |
复杂路由、多语言环境 |
| Apache Kafka |
自定义二进制协议 |
高吞吐、持久化、分区顺序 |
日志收集、流处理、大数据 |
| RocketMQ |
自定义协议 |
阿里出品,高吞吐,支持事务 |
电商、金融交易场景 |
2. JMS 与其他协议的关系
- JMS 是 API 规范(Java 接口标准),不是网络协议;
- AMQP 是网络协议(有线协议),RabbitMQ 是 AMQP 的典型实现;
- JMS 和 AMQP 是不同层面的概念——JMS 是 Java API 标准,AMQP 是跨语言的网络协议标准。
重要理解:
- 使用 JMS API 的应用,底层可以跑 AMQP 协议(如 ActiveMQ Artemis 同时支持 JMS 和 AMQP);
- JMS 的核心价值在于:无论底层用什么消息中间件,Java 应用代码都可以用同一套 JMS API 编写,具备可移植性。
3. JMS vs 主流 MQ 选型指南
| 场景 |
推荐方案 |
理由 |
| Jakarta EE 标准应用 |
ActiveMQ Artemis + JMS |
与 EJB、MDB 无缝集成,标准规范 |
| Spring Boot 项目 |
RabbitMQ 或 Kafka |
Spring 生态支持完善 |
| 高吞吐日志/流处理 |
Kafka |
百万级吞吐,持久化 |
| 复杂路由、多语言 |
RabbitMQ |
AMQP 协议,灵活路由 |
| 金融交易、强一致性 |
RocketMQ 或 ActiveMQ |
支持事务消息 |
| 老系统维护 |
ActiveMQ Classic |
稳定,JMS 规范 |
4. 核心选型原则
JMS 是 Java 世界的标准 API,但现代消息中间件大多已超越 JMS 范畴,提供了更丰富的特性。选型时应根据实际需求决定:需要标准化选 JMS 实现,需要高性能/大数据选 Kafka,需要灵活路由选 RabbitMQ。
九、总结
1. 核心知识脉络
| 层级 |
技术 |
核心概念 |
| 消息模型 |
Queue / Topic |
点对点(一对一)、发布订阅(一对多) |
| 核心 API |
Connection / Session / Producer / Consumer |
JMS 1.1 统一 API、JMS 2.0 简化 API |
| 消费方式 |
同步 / 异步 |
receive() 阻塞、MessageListener 回调 |
| 消息可靠性 |
持久化 / 事务 / 确认 |
PERSISTENT、本地事务、JTA/XA |
| 消息类型 |
TextMessage / MapMessage / 等 |
五种消息类型,TextMessage 最常用 |
2. 关键理解
- JMS 是 API 标准,不是产品:JMS 定义接口,具体实现由各消息中间件(ActiveMQ、Artemis 等)提供;
- 两种消息模型解决不同问题:Queue 解决任务分发(一对一),Topic 解决事件广播(一对多);
- 异步消费是生产环境首选:MessageListener 回调模式避免线程阻塞,提升系统吞吐量;
- 事务保证消息可靠性:本地事务保证单资源一致性,JTA/XA 保证多资源(JMS + 数据库)分布式一致性;
- JMS 与 EJB/MDB 天然集成:消息驱动 Bean(MDB)是 Jakarta EE 中消费 JMS 消息的标准方式(上一章已介绍)。
3. 与后续学习的衔接
理解 JMS 后,你将能更好地理解:
- Spring 的 JMS 支持:
JmsTemplate、@JmsListener 注解;
- 消息驱动架构:事件溯源(Event Sourcing)、CQRS;
- 微服务中的异步通信:服务间通过消息解耦,最终一致性。