背景

本文是《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 + 数据库)时,需要分布式事务保证一致性。

典型场景:订单服务收到请求后,需要同时:

  1. 向数据库写入订单记录;
  2. 发送 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. 业务需求

构建一个订单处理系统,包含以下流程:

  1. 用户下单 → 订单服务接收请求,保存订单到数据库;
  2. 订单服务发送 JMS 消息通知库存服务扣减库存;
  3. 库存服务异步消费消息,扣减库存;
  4. 扣减完成后,发送消息通知订单服务更新订单状态;
  5. 订单服务收到状态更新消息后,更新订单状态。

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. 关键理解

  1. JMS 是 API 标准,不是产品:JMS 定义接口,具体实现由各消息中间件(ActiveMQ、Artemis 等)提供;
  2. 两种消息模型解决不同问题:Queue 解决任务分发(一对一),Topic 解决事件广播(一对多);
  3. 异步消费是生产环境首选:MessageListener 回调模式避免线程阻塞,提升系统吞吐量;
  4. 事务保证消息可靠性:本地事务保证单资源一致性,JTA/XA 保证多资源(JMS + 数据库)分布式一致性;
  5. JMS 与 EJB/MDB 天然集成:消息驱动 Bean(MDB)是 Jakarta EE 中消费 JMS 消息的标准方式(上一章已介绍)。

3. 与后续学习的衔接

理解 JMS 后,你将能更好地理解:

  • Spring 的 JMS 支持JmsTemplate@JmsListener 注解;
  • 消息驱动架构:事件溯源(Event Sourcing)、CQRS;
  • 微服务中的异步通信:服务间通过消息解耦,最终一致性。