欢迎光临
我们一直在努力

Java-198 RabbitMQ JMS 模式详解:Queue/Topic、6 类消息与对象模型(JMS 2.0 / Jakarta Messaging 3.1)

TL;DR

  • 场景:Java 系统做异步解耦与事件驱动,需要统一理解 JMS 的消息模型、对象模型与消息类型。
  • 结论:JMS 是标准 API(类似 JDBC),关键在 Queue/Topic 语义、Session 与确认/事务边界、消息类型取舍。
  • 产出:一篇可直接落地的 JMS 基础稿:概念—对象模型—消息类型—实现/版本差异速查。

请添加图片描述

版本矩阵

项目说明
JMS 2.0 属于 Java EE 7,API 仍在 javax.jms 命名空间,可在 Java SE 环境独立使用。
Jakarta Messaging 3.1 属于 Jakarta EE 10,对应 API 命名空间为 jakarta.jms(主要是包名变更)。
ActiveMQ Classic 完整支持 JMS 1.1,对 JMS 2.0 / Jakarta Messaging 3.1 属于“部分支持”,强调 jakarta.jms 迁移价值。
ActiveMQ Artemis JMS/Jakarta Messaging 标准不定义网络协议;Artemis 客户端基于自身协议实现,并提供客户端侧 JNDI。
RabbitMQ 本身不是 JMS Provider,需要配套插件与 RabbitMQ JMS Client 才能以 JMS Queue/Topic 方式对接。

JMS模式

基本介绍

JMS(Java Message Service)是Java平台提供的一套标准API,专门用于实现面向消息中间件(MOM,Message Oriented Middleware)的编程接口。作为Java EE规范的重要组成部分,JMS定义了一组通用接口和语义,允许Java应用程序通过消息队列或发布/订阅模式进行异步通信。

核心特点:

  • 跨平台性:JMS提供与具体实现无关的抽象接口,使得应用程序可以与任何兼容JMS规范的MOM提供商(如ActiveMQ、RabbitMQ、IBM MQ等)交互
  • 松耦合通信:通过消息代理(Message Broker)实现生产者和消费者的解耦,双方无需同时在线或直接交互
  • 支持两种消息模型:
    • 点对点(Point-to-Point):基于队列(Queue)的精确一次消费模式
    • 发布/订阅(Publish/Subscribe):基于主题(Topic)的多消费者广播模式
  • 技术类比: JMS与JDBC(Java Database Connectivity)在架构设计理念上高度相似:

    • 都采用"接口+实现"的设计模式
    • 都定义标准API,由不同厂商提供具体实现
    • 都通过配置切换底层服务提供商而不需修改应用代码

    典型应用场景:

  • 异步任务处理(如订单处理、日志收集)
  • 系统解耦(微服务架构中的服务间通信)
  • 事件驱动架构(EDA)的实现基础
  • 跨系统数据同步(如库存系统与物流系统)
  • 通过JMS,Java开发者可以构建高可靠、可扩展的分布式消息系统,有效解决传统同步通信带来的性能瓶颈和系统耦合问题。

    JMS消息

    消息是JMS中的一种类型对象,由两部分组成:报文头和消息主题。 报文头包括消息头字段和消息头属性,字段是JMS协议规定的字段,属性可以由用户按需添加。 JMS 报文头全部字段:

    在这里插入图片描述

    消息主题携带应用程序的数据或有效负载,根据有效负载的不同格式和用途,JMS(Java Message Service)规范定义了以下几种标准消息类型,每种类型都有其特定的应用场景和优势:

  • 简单文本(TextMessage)

    • 最常用的消息类型,包含普通字符串内容
    • 示例:JSON/XML格式的业务数据、简单的通知文本
    • 应用场景:订单信息、日志消息、系统通知等
  • 可序列化对象(ObjectMessage)

    • 包含一个可序列化的Java对象
    • 要求传输的对象必须实现java.io.Serializable接口
    • 示例:传输完整的用户对象、产品对象等业务实体
    • 注意:跨系统使用时需确保双方有相同的类定义
  • 属性集合(MapMessage)

    • 包含一组键值对数据,键为String类型,值为Java基本类型
    • 类似Properties的概念,但支持更多数据类型
    • 示例:配置参数、表单数据、属性集合
  • 字节流(BytesMessage)

    • 包含未解释的字节流数据
    • 适用于传输二进制数据
    • 示例:文件传输、图片/视频数据、加密数据包
    • 开发者需要自行处理字节序和编解码
  • 原始值流(StreamMessage)

    • 包含一系列Java基本类型的流数据
    • 类似数据流的顺序读写操作
    • 示例:传感器数据流、实时监控数据
  • 无有效负载的消息(Message)

    • 不包含实际数据,仅包含消息头和属性
    • 用于事件通知或触发操作
    • 示例:系统心跳、事件触发器、简单的通知信号
  • 每种消息类型都支持设置消息属性(Properties),这些属性可以包含消息的元数据,如:

    • 消息优先级
    • 过期时间
    • 持久性设置
    • 自定义业务属性

    在实际应用中,选择哪种消息类型取决于:

  • 数据的性质和结构
  • 系统间的兼容性需求
  • 性能考虑(如序列化开销)
  • 消息处理逻辑的复杂度
  • 体系架构

    JMS由以下元素组成:

    • JMS 供应商产品:JMS接口的一个实现,该产品可以是Java的JMS实现,也可以是非Java的面向消息中间件的适配器。
    • JMS Client:生产或者消费基于消息的Java的应用程序或者对象
    • JMS Producer:创建并发送消息的JMS客户
    • JMS Consumer:接收消息的 JMS 客户
    • JMS Message:包括可以在 JMS 客户之间传递的数据的对象
    • JMS Queue:缓存消息的容器,消息的接受顺序并不一定非要与消息的发送顺序相同,消息被消费后将从消息队列中移除
    • JMS Topic:Pub、Sub模式

    对象模型

    在这里插入图片描述

    Connection Factory

    Connection Factory(连接工厂)是JMS(Java Message Service)中一个重要的被管对象,它主要负责创建和管理与JMS提供者之间的连接。作为JMS架构中的关键组件,连接工厂提供了以下核心功能:

  • 连接创建机制

    • 通过统一的接口创建Connection对象
    • 封装了底层连接建立的复杂细节
    • 支持不同类型的连接方式(TCP/IP、HTTP等)
  • 配置管理

    • 管理员通过JNDI(Java Naming and Directory Interface)名字空间配置连接工厂
    • 可以设置连接参数(超时时间、重试次数等)
    • 支持集群配置和高可用性设置
  • 类型区分

    • 根据消息传递模型的不同分为两种类型:
      • QueueConnectionFactory(队列连接工厂):用于点对点消息模型
      • TopicConnectionFactory(主题连接工厂):用于发布/订阅消息模型
  • 可移植性优势

    • 提供标准的JMS接口,与具体实现解耦
    • 当更换JMS提供商时,客户端代码无需修改
    • 支持多种JMS实现(ActiveMQ、RabbitMQ等)
  • 应用场景示例

    • 企业应用中,通过连接工厂创建到消息服务器的连接
    • 分布式系统中,使用连接工厂实现跨系统的异步通信
    • 云环境中,通过配置不同的连接工厂适应不同的部署环境
  • 使用流程

  • // 1. 通过JNDI查找连接工厂
    Context ctx = new InitialContext();
    ConnectionFactory cf = (ConnectionFactory)ctx.lookup("ConnectionFactory");

    // 2. 创建连接
    Connection connection = cf.createConnection();

    // 3. 创建会话
    Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);

    // 4. 使用消息
    // …

  • 性能考量
    • 连接工厂通常配置连接池提高性能
    • 可以设置最大连接数等参数
    • 支持SSL等安全连接方式
  • 通过连接工厂的这种设计,JMS实现了客户端代码与具体消息服务的解耦,提高了系统的灵活性和可维护性。

    Connection

    连接,连接代表了应用程序和消息服务器之间的通信链路,在获得了连接工厂之后,就可以创建一个与JMS提供者的连接,根据不同的连接类型,连接允许用户创建会话,以发送和接收队列主题为目标。

    Destination

    目标是一个包装了消息目标标识符号的被管对象,消息目标是指消息发布和接收的地点,或者是队列,或者是主题,JMS管理员创建这些对象,然后用户通过JNDI发现他们,和连接工厂一样,管理员可以创建两种管理类型的目标,点对点模型的队列,以及发布者、订阅者模型的主题。

    Session

    表示一个单线程的上下文,用于发送和接收消息,由于会话是单线程的,所以消息是连续的,就是说消息是按照顺序一个一个接收的。 会话好处是它支持事务,如果用选择了事务支持,会话上下文会将保存一组消息,直到事务被提交了才发送这些消息。 在提交事务之前,用户可以使用回滚操作取消这些消息,一个会话允许用户创建消息,生产者发送消息,消费者来接收消息。

    MessageConsumer

    MessageConsumer 是 JMS (Java Message Service) 中的核心接口之一,代表消息消费者。它由 Session 会话对象创建,专门用于接收发送到 Destination(目标)的消息。消息消费者可以接收两种类型的消息目标:

  • 队列(Queue):点对点消息模型,每条消息只能被一个消费者接收
  • 主题(Topic):发布/订阅消息模型,消息会被广播给所有订阅者
  • MessageConsumer 提供两种消息接收模式:

    同步接收(阻塞模式):

    • 通过 receive() 方法实现
    • 线程会一直阻塞直到收到消息或超时
    • 示例代码:

    Message message = consumer.receive(); // 无限等待
    Message message = consumer.receive(5000); // 等待5秒

    异步接收(非阻塞模式):

    • 通过设置 MessageListener 实现
    • 采用事件驱动机制,当消息到达时会自动触发 onMessage 方法
    • 示例代码:

    consumer.setMessageListener(new MessageListener() {
    @Override
    public void onMessage(Message message) {
    // 处理接收到的消息
    }
    });

    重要特性:

    • 支持消息选择器(MessageSelector)来过滤消息
    • 可以配置消息确认模式(AUTO_ACKNOWLEDGE, CLIENT_ACKNOWLEDGE等)
    • 支持持久化订阅(Durable Subscription)用于主题模式

    MessageProducer

    MessageProducer 是 JMS 中用于发送消息的核心接口,由 Session 会话对象创建。它负责将消息发送到指定的 Destination(目标)。MessageProducer 有两种创建方式:

  • 特定目标生产者:
  • MessageProducer producer = session.createProducer(destination);

  • 通用生产者:
  • MessageProducer producer = session.createProducer(null);
    // 发送时指定目标
    producer.send(specificDestination, message);

    主要功能特性:

    消息发送控制:

    • 可设置消息的存活时间(TimeToLive)
    • 可配置消息优先级(Priority)
    • 支持持久化/非持久化消息传递模式

    发送方法:

    // 简单发送
    producer.send(message);

    // 带参数发送
    producer.send(message, deliveryMode, priority, timeToLive);

    典型应用场景:

  • 订单处理系统:生产者发送订单消息到队列
  • 新闻推送系统:生产者发布新闻到主题
  • 日志收集系统:多个生产者发送日志消息到集中队列
  • 注意事项:

    • 发送完成后应调用 close() 方法释放资源
    • 可通过 setDeliveryDelay() 设置消息延迟投递
    • 支持事务性发送(当使用事务会话时)

    Message

    消息是分布式系统中用于在生产者(Producer)和消费者(Consumer)之间传输数据的核心载体。它本质上是一个结构化的数据对象,用于在不同应用程序或服务间实现异步通信和解耦。消息队列(Message Queue)或消息中间件(如RabbitMQ、Kafka)通常作为消息的传输媒介,确保消息的可靠传递和顺序性。

    一个完整的消息通常由以下三部分组成:

  • 消息头(Header,必须) 消息头是消息的元数据部分,包含用于控制消息路由和处理的关键信息。常见的消息头字段包括:

    • Message ID:唯一标识符,用于追踪消息。
    • Timestamp:消息的创建或发送时间。
    • Priority:消息的优先级(如高、中、低),用于决定处理顺序。
    • Expiration:消息的有效期,超时后可能被自动丢弃。
    • Destination:目标队列或主题的名称,用于路由。
    • Correlation ID:关联ID,常用于请求-响应模式中匹配请求和回复。

    示例:在订单处理系统中,消息头可能包含订单ID作为Correlation ID,确保订单状态更新能正确关联到原始请求。

  • 消息属性(Properties,可选) 消息属性是一组可扩展的键值对,用于补充消息的附加信息或实现特定功能。常见的用途包括:

    • 自定义过滤:通过属性定义标签(如region=Asia),消费者可使用消息选择器(Message Selector)订阅特定消息。
    • 兼容性扩展:支持不同消息中间件的特性(如RabbitMQ的headers交换器或Kafka的消息标签)。
    • 业务上下文:添加业务相关的字段(如payment_method=credit_card)。

    示例:在物流系统中,消息属性可能包含shipping_priority=express,供消费者快速处理加急订单。

  • 消息体(Body,可选) 消息体是消息的实际负载数据,支持多种格式以适应不同场景。常见的消息体类型包括:

    • 文本消息(TextMessage):纯文本内容,如JSON或XML格式的订单数据。
    • 映射消息(MapMessage):键值对结构,适合传递动态字段(如{"userId": 123, "status": "shipped"})。
    • 字节消息(BytesMessage):二进制数据,适用于文件或图像传输。
    • 流消息(StreamMessage):类似数据流的连续字节序列,通常用于实时处理。
    • 对象消息(ObjectMessage):序列化的Java对象(需注意跨语言兼容性问题)。

    示例:在电商系统中,订单创建消息可能使用JSON格式的文本消息体,包含商品列表、用户地址等信息。

  • 应用场景

    • 异步处理:订单系统将支付成功消息发送到队列,库存服务异步消费并扣减库存。
    • 事件驱动架构:用户注册事件触发邮件服务和数据分析服务的并行处理。
    • 跨系统集成:银行系统通过消息属性标记交易类型(如txn_type=transfer),便于风控系统过滤处理。

    错误速查

    症状根因定位修复
    NoClassDefFoundError: javax/jms/… 或 jakarta/jms/… 依赖命名空间不匹配(javax.jms vs jakarta.jms),或缺失 API 包 看异常类名属于 javax 还是 jakarta;查依赖树是否同时混入两套 API 统一到一套 API;Spring 6/ Jakarta EE 9+ 生态通常要求 jakarta.jms(注意与 Broker 客户端兼容)
    Topic 订阅收不到历史消息 非持久订阅(non-durable)只接收订阅后消息;或订阅未正确建立 检查是否配置 durable subscription、clientId、subscriptionName;确认是否使用 Topic 而非 Queue 需要历史消息就用 durable subscription;同时确保 clientId 与订阅名稳定(重启不变)
    消费者一直阻塞 / 收不到消息 未调用 connection.start();或使用了错误 Destination;或 Selector 过滤掉 打日志确认 start 调用顺序;打印 Destination 名称;临时去掉 selector 按“创建连接→创建 session→创建 consumer→start→receive/onMessage”顺序;先无 selector 验证链路
    重复消费 / “处理过但又来一遍” ACK/事务边界不正确;处理异常触发重投;超时导致回滚 看确认模式(AUTO/CLIENT/DUPS_OK/事务);看 broker redelivery 日志/死信队列 要么开启事务并在成功后 commit;要么 CLIENT_ACKNOWLEDGE 并在业务成功后 ack;失败明确回滚并配 DLQ 策略
    消息顺序不稳定 多消费者并发、预取、分区/分片策略导致乱序 同一队列是否有多个 consumer;是否启用并发 listener;看 broker 的 dispatch/consumer 数量 需要强顺序:单 consumer(或按业务 key 分队列/分区);避免在同一队列上盲目扩并发
    ObjectMessage 反序列化失败 / ClassNotFoundException 消费端缺少同名类或 serialVersionUID 不兼容;跨语言天然不适配 看异常类名;对比生产/消费端 jar 版本;检查是否跨服务/跨团队 ObjectMessage 仅限同构 Java 且版本强约束;跨系统用 TextMessage(JSON) 或 BytesMessage 自定义协议
    RabbitMQ 上 JMS 行为异常(Topic/Selector 不符合预期) RabbitMQ 不是原生 JMS Provider,行为依赖插件与 JMS Client 映射层 确认是否启用对应插件;核对 JMS Client/插件版本是否匹配 按官方说明启用插件并使用 RabbitMQ JMS Client;不要假设其语义与 ActiveMQ/IBM MQ 完全一致
    迁移到 jakarta.jms 后“能编译但连不上/不兼容” API 包名升级≠协议兼容;broker/客户端组合不支持该迁移路径 核对 broker 类型(Classic/Artemis/其他)与客户端实现;看官方兼容说明 先确认目标 broker 是否有对应 jakarta.jms 客户端链路;必要时升级 broker 或改用支持 Jakarta 的实现

    其他系列

    🚀 AI篇持续更新中(长期更新)

    AI炼丹日志-29 – 字节跳动 DeerFlow 深度研究框斜体样式架 私有部署 测试上手 架构研究,持续打造实用AI工具指南! AI研究-132 Java 生态前沿 2025:Spring、Quarkus、GraalVM、CRaC 与云原生落地 🔗 AI模块直达链接

    💻 Java篇持续更新中(长期更新)

    Java-196 消息队列选型:RabbitMQ vs RocketMQ vs Kafka MyBatis 已完结,Spring 已完结,Nginx已完结,Tomcat已完结,分布式服务已完结,Dubbo已完结,MySQL已完结,MongoDB已完结,Neo4j已完结,FastDFS 已完结,OSS已完结,GuavaCache已完结,EVCache已完结,RabbitMQ正在更新… 深入浅出助你打牢基础! 🔗 Java模块直达链接

    📊 大数据板块已完成多项干货更新(300篇):

    包括 Hadoop、Hive、Kafka、Flink、ClickHouse、Elasticsearch 等二十余项核心组件,覆盖离线+实时数仓全栈! 大数据-278 Spark MLib – 基础介绍 机器学习算法 梯度提升树 GBDT案例 详解 🔗 大数据模块直达链接

    赞(0)
    未经允许不得转载:171主机测评 » Java-198 RabbitMQ JMS 模式详解:Queue/Topic、6 类消息与对象模型(JMS 2.0 / Jakarta Messaging 3.1)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址