
做消息中间件开发这些年我见过最吊诡的一幕同一个团队写着差不多的对接逻辑对着ActiveMQ能用换成RabbitMQ就一串诡异报错明明代码里调用了Session.commit()消息却还是没被消费。追究到最后所有人的目光都转向了那四个字——JMS标准。JMS不是某一个软件更不是ActiveMQ的专属名词。它是Java领域里最早成体系的消息服务规范定义了消息中间件在Java平台上的统一访问方式。很多人觉得JMS已经过时毕竟现在Kafka、Pulsar这些新东西满天飞但如果你去翻Spring的源码会发JmsTemplate至今仍是Spring消息模块的基石。搞懂JMS标准不只是为了继承传统代码更是为了理解标准这件事在中间件领域到底是怎么运作的——它规定了什么放开了什么哪些坑是规范本身造成的哪些坑是实现厂商的锅。这篇文章我会从JMS规范的核心模型讲起结合真实代码和踩坑经历把消息头、可靠性机制、事务、订阅模式、与AMQP/Kafka的差异这些事说透。适合正在用或准备用JMS API的Java工程师以及想搞清楚JMS标准和ActiveMQ文档有什么区别的技术人。1. 我们真的需要JMS标准吗先聊聊它解决过什么1.1 从一次消息重复消费事故说起有一回我排查线上问题业务方反馈订单支付成功后返券消息被消费了两次导致用户收到两张券。查来查去代码里用了Session.AUTO_ACKNOWLEDGE消费者在监听方法里查了数据库但异常发生在消息回调返回之后、确认信号发出之前。这时消息会被重新投递而消费者的本地事务已经提交了就造成重复处理。这不是框架的bug是我对JMS可靠性模型的理解不够。JMS标准把消息是否已处理交给会话确认但处理的边界由应用自己定义。AUTO_ACKNOWLEDGE模式下确认发生在onMessage正常返回后如果消息回调中途抛出异常或这个返回和确认之间存在崩溃窗口消息就可能被重投。标准文档里写得很清楚它保证了至少一次语义不保证恰好一次。这个特性同样适用于大多数消息中间件理解了JMS的措辞你就不会对Kafka的高水位或RabbitMQ的手工ack感到陌生。1.2 JMS在Java技术栈里的生态位置在JMS出现之前Java应用要与不同消息中间件对接只能使用各家私有API。如果今天用IBM MQ明天换成Oracle AQ业务代码就废了。JMS标准做的事情就是定义一组统一的接口——ConnectionFactory、Connection、Session、Destination、MessageProducer、MessageConsumer、Message。只要中间件提供了JMS实现你就能用同一套代码跑在不同的消息系统上。它解决的是Java应用与消息中间件之间的标准化接入问题而不是消息如何路由的问题。路由、存储、重试策略这些都是中间件自身的事JMS只约束客户端视角的行为。这一点跟JDBC很像JDBC规定了Java怎么连数据库、怎么执行SQL但MySQL和PostgreSQL的存储引擎差异JDBD管不着。1.3 为什么今天仍然要学JMS标准有些人觉得JMS过时了因为Kafka用的是自己的一套协议RocketMQ也提供了更丰富的特性。但市面上仍有大量企业级应用基于ActiveMQ、Artemis、WebLogic JMS跑着Spring的JmsTemplate和JmsListener也是JMS API的封装。你直接调用JMS接口写程序换一个兼容JMS的Broker代码几乎不用改这就是标准的价值。而且JMS标准中包含了很多设计理念比如会话、消息头、持久订阅、事务边界。这些概念不是ActiveMQ发明的是JMS规范抽象出来的。你学会了JMS再看AMQP的Basic.Consume和channel.basicAck会有一种强烈的既视感——本质上不同标准是在用不同的语言描述同一件事。2. 拆开JMS标准的地图五个核心接口撑起一个消息世界2.1 ConnectionFactory与Connection从拨号到上线JMS的客户端启动过程是有层次的先拿到连接工厂再创建连接。连接工厂在JNDI里配置携带Broker的地址、认证信息。创建连接是一个比较重的操作它要建立TCP连接、握手、可能还要认证授权所以不要在每次发送消息时都新建连接。ConnectionFactory factory new ActiveMQConnectionFactory(tcp://localhost:61616); Connection connection factory.createConnection(admin, admin); connection.start();很多新手漏掉connection.start()以为创建了连接就能收消息。结果就是生产者正常发送消费者监听回调却一次也不触发。规范里的设计是连接建立后默认处于停止状态要显式调用start()才能开始投递消息。这在连接复用、先加监听器再放量的场景下非常合理但你如果不知道这个点会浪费一个下午。连接是线程安全的多个线程可以共享一个Connection。但Connection内部包含的Session不是线程安全的这就要说下一个核心接口。2.2 Session你的一切消息操作都在这里发生Session是JMS里最重要的边界。它既是消息生产者、消费者的工厂也是事务和确认机制的工作单元。每次创建会话时都要指定事务模式和确认模式Session session connection.createSession(false, Session.AUTO_ACKNOWLEDGE);第一个参数transacted如果为true第二个参数会被忽略会话运行在事务模式。第二个参数决定非事务模式下的消息确认策略。Session不是线程安全的推荐的使用方式是一个线程一个Session或者用池技术复用。如果多个线程共享一个Session经常会出现消息丢失、消费错乱、死锁等问题而且极难排查。从JMS 2.0开始规范引入了JMSContext一个JMSContext内部管理一个Connection和一个Session并通过try-with-resources自动关闭简化了老的连接/会话双层模型。但老代码存量太大你还是得认识Session。2.3 Destination与MessageProducer/MessageConsumer发到哪儿谁去取Destination是队列和主题的统一抽象。Queue是点对点模型一条消息只有一个消费者Topic是发布订阅模型一条消息能广播给所有订阅者。生产者和消费者都是基于Destination创建的Queue queue session.createQueue(ORDER.QUEUE); MessageProducer producer session.createProducer(queue); TextMessage msg session.createTextMessage(订单创建成功); producer.send(msg); MessageConsumer consumer session.createConsumer(queue); Message received consumer.receive(5000);这里有个容易忽略的细节session.createProducer(queue)和producer.send(msg)还有另一个重载可以在send时指定Destination。也就是说生产者可以与某个Destination绑定也可以在每条消息上动态指定目的地。后者更灵活但容易造成消息发错队列的生产事故不建议随意使用。消费者有两种接收方式同步的receive()和异步的setMessageListener()。同步阻塞时如果Broker断开且没有设置较长的receiveTimeout可能卡很久。异步监听器更适合高吞吐场景但在多线程消费时要注意监听器单线程模型的限制。3. 消息本身的仪式感往返头、属性与Body的约定3.1 JMS头字段JMSMessageID、JMSExpiration的真正用途JMS消息由三部分构成消息头、属性、消息体。头字段是协议级的元数据比如JMSDestination消息实际到达的目的地、JMSDeliveryMode持久或非持久、JMSMessageID全局唯一ID、JMSTimestamp发送时间、JMSExpiration过期时间、JMSPriority优先级。其中JMSMessageID是去重的最佳切入点。但要注意JMS规范并没有要求消息ID绝对唯一只是说生产者应生成一个唯一值。不同实现可能用不同算法ActiveMQ生成的ID形如ID:hostname-xxx你如果依赖这个字段做幂等键得确认Broker的生成策略是否可靠。JMSExpiration也常被忽略。默认值为0表示永不过期但消息中间件往往有全局的TTL配置。你发送时没有设置过期时间Broker却可能按自己的默认TTL把消息丢了。我踩过的坑是ActiveMQ的messageTTL默认值在某些版本里是0不过期在Artemis里默认是0毫秒立即可过期不对Artemis默认是永不过期但有些云服务会把TTL设成3天。生产环境一定要显式设置setTimeToLive别让不规范实操影响到线上体验。3.2 应用属性和选择性消费用选择器把数据过滤在到达之前JMS允许在消息上挂载自定义属性属性值可以是boolean、byte、short、int、long、float、double、String这些类型。消费者可以使用选择器表达式来过滤消息MessageConsumer consumer session.createConsumer(queue, orderType vip AND amount 100);选择器类似于SQL的WHERE子句但要注意语法差异字符串字面量必须用单引号属性名区分大小写IS NULL和IS NOT NULL可以判断属性是否存在。选择器的运算是在Broker端完成的吗不完全对。JMS规范允许客户端或Broker实现过滤但很多Broker是在服务端过滤的。即便如此你在消费端写选择器时还是要尽量利用消息头属性因为服务端过滤是对用户透明的优化客户端也可能收到一些不满足条件的消息其实不会但和Broker实现有关。这个机制最大的价值是把过滤逻辑前移减少无效消息的传输和处理。不过也有代价选择器在Broker端会增加索引和匹配计算如果大量消费者使用复杂选择器Broker CPU可能升高。一般只建议对高筛选价值属性建选择器别把选择器当SQL join使。3.3 Message Body的五种类型以及何时用BytesMessage/MapMessageJMS定义了五种消息体TextMessage文本字符串最常见。MapMessage键值对集合。BytesMessage字节数组适合二进制数据。StreamMessage顺序读取的原始值流。ObjectMessage可序列化Java对象。ObjectMessage看起来方便但我强烈建议少用。它的序列化形式绑定了Java类版本如果发送方和接收方类定义不一致反序列化直接报ClassNotFoundException。系统升级时新旧版本消息混在队列里很容易炸。如果一定要用建议在消息里带上classVersion属性并做好兼容处理。BytesMessage则是最通用的方式。不同语言的消息中间件客户端比如C#、Python都能解析字节流所以跨语言场景优先用BytesMessage。MapMessage也还行但Map的键值在跨语言时语义可能不一致不如直接定义一个明确的二进制协议。4. 可靠性不是一句保证不丢确认、持久化与订阅机制4.1 AUTO_ACKNOWLEDGE到SESSION_TRANSACTED四种确认模式选谁JMS规范定义了四种确认模式它们决定了消息何时被标记为已消费。模式行为适用场景AUTO_ACKNOWLEDGE消息回调正常返回后自动确认对重复不敏感、处理逻辑简单的场景CLIENT_ACKNOWLEDGE应用显式调用message.acknowledge()需要批量确认或处理后有后续业务动作DUPS_OK_ACKNOWLEDGE延迟确认允许重复消息高吞吐、可容忍重复的非关键场景SESSION_TRANSACTED事务提交时统一确认需要事务保证多消息/数据库原子性CLIENT_ACKNOWLEDGE有个非常容易踩坑的细节Message.acknowledge()会确认当前消息之前所有未确认的消息而不是只确认这一条。这是因为JMS的确认是基于Session的相当于确认游标往前推进。很多刚接触的人以为调了message.acknowledge()就只确认这条结果把之前积压的消息全确认掉了造成数据丢失。DUPS_OK_ACKNOWLEDGE模式下Broker可以延迟确认消息允许生产者重发。这个模式很少用但如果你在做日志采集这类可乱序、可重复的场景它确实能减少确认开销。4.2 持久订阅与非持久订阅以及Durable Subscription的坑点对点模型下消息会保存在Queue里消费者不在时消息不会丢取决于持久化配置。但发布订阅模型下如果一个订阅者不在线非持久订阅的Topic消息会直接被丢弃。JMS为解决这个问题定义了持久订阅Durable Subscription。创建持久订阅需要指定客户端IDclientID并且订阅名称要唯一Connection connection factory.createConnection(); connection.setClientID(order-consumer-01); Session session connection.createSession(false, Session.AUTO_ACKNOWLEDGE); Topic topic session.createTopic(ORDER.TOPIC); MessageConsumer consumer session.createDurableSubscriber(topic, order-subscriber);这里有几个坑第一clientID在一个连接内是唯一的如果你在多个连接上使用同一个clientIDBroker会直接抛错。第二持久订阅会保存订阅状态即使消费者关了消息也会被保留重新连接后会补发。但如果应用不再需要这个订阅必须调用unsubscribe清理否则消息会在Broker堆积直到击穿内存。第三持久订阅和普通消费者不同它本质上是在Broker上为这个订阅建了一个独立的队列消费逻辑要自己维护去重和游标。4.3 事务让收发一体成为原子操作JMS事务把一个Session内所有生产者和消费者的操作统一包裹在一个原子单元里。在事务会话中Session session connection.createSession(true, Session.SESSION_TRANSACTED); try { producer.send(msg1); producer.send(msg2); session.commit(); } catch (JMSException e) { session.rollback(); }commit()时会一次性确认所有已接收消息并批量发送所有待发送消息rollback()会取消本事务内所有操作。注意事务会话只能使用MessageProducer和MessageConsumer的同步/异步操作在同一个Session中完成如果你用多个线程操作同一个事务Session会引发不可预期行为。还有一个概念是XA事务。JMS规范支持XASession让消息操作和数据库操作处于同一个全局事务中。比如订单表和消息发送同时成功、同时失败。但XA事务的开销非常大而且一旦Broker或数据库一方超时整个事务卡死排查起来堪比破案。我的建议是能用本地消息表定时补偿就别上XA。用数据库事务先写业务表和消息表再由一个定时任务把消息表里未发送的记录捞出来投递是更可控的方案。5. 代码落地从ActiveMQ到原生JMS的完整链路5.1 搭建JMS依赖不要跳过Broker配置细节以ActiveMQ为例Maven依赖只需要dependency groupIdorg.apache.activemq/groupId artifactIdactivemq-client/artifactId version5.18.4/version /dependency但真正常被忽略的是Broker端配置内存限制、消息持久化策略、连接数限制。如果本地运行ActiveMQ出现java.lang.IllegalStateException: Broker not started大多是Broker启动失败或端口被占用。用Docker起一个测试实例最省心docker run -p 61616:61616 -p 8161:8161 -e ACTIVEMQ_ADMIN_LOGINadmin -e ACTIVEMQ_ADMIN_PASSWORDadmin -v activemq_data:/data -v activemq_conf:/conf --name activemq rmohr/activemq:5.18.4启动后管理控制台在http://localhost:8161暗坑是很多教程不告诉你ActiveMQ默认的openwire协议端口61616和管理端口8161连接失败时先检查端口。5.2 生产者发送优先用事务批处理和持久化开关生产者最简单实现try (JMSContext context new ActiveMQConnectionFactory(tcp://localhost:61616).createContext(admin, admin)) { JMSProducer producer context.createProducer(); producer.setDeliveryMode(DeliveryMode.PERSISTENT); producer.setTimeToLive(60000); producer.send(context.createQueue(ORDER.QUEUE), new order from user#123); }注意我用的是JMS 2.0的JMSContext比老API简洁很多。setDeliveryMode如果设置成非持久模式NON_PERSISTENTBroker宕机后未落盘的消息会全部丢失。很多内部日志场景可以接受但交易类消息必须要用PERSISTENT。对于大批量发送建议把多条消息放在一个事务Session里批量提交性能远优于逐条发送。ActiveMQ使用单个事务批量发送吞吐量能提升一个数量级。JMS本身并不禁止在循环里Send但多次网络往返的开销肉眼可见。5.3 消费者落地监听器 手动确认的正确写法建议用CLIENT_ACKNOWLEDGE包一层业务逻辑Connection conn factory.createConnection(); conn.setClientID(consumer-1); Session session conn.createSession(false, Session.CLIENT_ACKNOWLEDGE); Queue queue session.createQueue(ORDER.QUEUE); MessageConsumer consumer session.createConsumer(queue); consumer.setMessageListener(message - { try { processMessage(message); // 业务处理 message.acknowledge(); // 明确告知Broker成功 } catch (Exception e) { // 记录日志进入重试或死信队列不要立即ack } }); conn.start();这里有一个很关键的细节conn.start()要放在setMessageListener之后。虽然JMS没有强制但你先start了再设置监听器很可能丢失在设置瞬间到达的消息。另外在CLIENT_ACKNOWLEDGE模式中如果处理失败了但没有ack消息会一直待在Broker里重启后会重投。有人会问那我想把失败的消息立即移动到死信队列怎么办JMS规范没定义死信队列那是Broker的扩展功能。ActiveMQ的RedeliveryPolicy和死信队列配置是另一套体系但基础的失败重投逻辑还是基于确认机制的。6. 我看过的坑连接泄漏、确认模式误用、异步消费性能陷阱6.1 连接和会话泄漏的排查最常见的故障不是JMS功能问题而是资源泄漏。很多人在每次发送时创建Connection和Session用完不关闭最终Broker连接数爆掉客户端报javax.jms.JMSException: Could not connect to broker。正确做法是复用连接或者用连接池。ActiveMQ官方推荐PooledConnectionFactorySpring的CachingConnectionFactory也能缓存Session。手工管理连接时记得在finally里关闭Session和Connection并且close()的顺序是先关Session再关Connection否则可能出现资源未释放的警告。如果你发现系统运行几天后连接数持续上涨先用lsof -p pid | grep 61616看连接数再用JVisualVM看Connection对象的数量很快就能定位到哪个模块没有close。6.2 事务回调放在错误线程异步监听器配合事务Session时有个隐蔽的问题。onMessage回调运行在JMS消息分发线程而业务代码里如果用了线程池异步处理主流程可能在onMessage返回时就认为处理成功但真正的业务处理还没完成。这导致事务提交了但业务数据还没写库或者反过来业务处理失败但消息已经确认。我的建议是事务边界必须与消息确认边界保持一致。也就是说在onMessage里同步完成所有需要和消息确认同生死的操作再返回。如果想异步处理就用手动ack在异步任务最后调用message.acknowledge()但要注意消息确认的对象还是同一个Session跨线程调用Session是不安全的。这种情况下更好的方案是使用队列批量拉取自己做线程池调度。6.3 大批量消费时选择器与预取的权衡JMS规范的createConsumer可以选择器这很强大但选择器要是写得太复杂Broker端扫描所有消息匹配属性性能会直线下降。另外同一个Connection上的多个消费者会共享网络I/O如果其中一个消费者选择器能匹配海量消息可能把其他消费者的消息抢走。ActiveMQ的prefetch limit预取限制也是一个变量ActiveMQPrefetchPolicy默认值对不同消费模型不同如果设置过大一个消费者可能一次性拉走大量消息别的消费者分不到。解决这类问题一方面要控制选择器的属性基数避免在非索引属性上做逻辑复杂的运算另一方面要根据机器的处理能力设置合理的prefetch。比如消费端逻辑很重时把prefetch设成1让消息一条一条处理反而能降低内存压力和重试成本。7. JMS不是唯一标准它和AMQP、MQTT、Kafka的定位差异很多人以为JMS标准等同于ActiveMQ或者觉得JMS和AMQP是对立关系。其实它们是不同层面的东西。JMS是Java API层面的标准AMQP是协议层面的标准MQTT是面向受限物联网的轻量协议Kafka则是自带一套自定义协议的分布式日志系统。标准/协议所属层面核心消息模型主要使用场景JMSJava API规范Queue/TopicJava应用间ActiveMQ、ArtemisAMQP线级协议Exchange/Queue/Binding跨语言RabbitMQMQTT发布/订阅协议Topic QoS物联网、低带宽Kafka自定义协议Partition Consumer Group大数据流、事件溯源JMS客户端与Broker之间交互的语言JMS规范并没有强行规定可以是OpenWire可以是STOMP甚至可以是AMQP映射。所以ActiveMQ支持JMS和ActiveMQ支持AMQP并不冲突。前者是给Java程序员提供了一套标准化API后者是给多语言客户端提供了统一的线路协议。从这个角度讲JMS标准最大的价值是它锁住了应用代码让业务不绑定到某个厂商的具体API。如果你有一天从ActiveMQ换到Artemis只要后者实现了JMS规范业务代码基本不用改。这比你在代码里到处用RocketMQClient的私有API要省心得多。当然这也意味着你无法轻易使用到中间件特有的高级特性比如ActiveMQ的延迟投递、消息组这些JMS规范都不包含。所以实际落地时标准API加厂商扩展API的混合使用是很多项目不得不走的路。8. 学习JMS标准应该抓住的主线如果让我给刚入行的同事列一个JMS学习路径我不会让他背接口而是从三条主线走消息模型、可靠性语义、事务边界。这三条线研究透了再看任何具体中间件都很快。第一消息模型。Queue和Topic的本质区别是什么持久订阅到底发生了什么。把这些画成图用ActiveMQ管理后台看消息流动比光看文档有效得多。第二可靠性语义。确认模式、持久化、重投、幂等它们如何组合出不丢不重的效果。重要的是理解不重在很多场景下是应用层的责任不是JMS标准能完全保证的。第三事务边界。JMS事务与数据库事务的关系何时该用事务会话何时该用外部补偿。这个分寸感体现的是架构能力。我见过太多人一上来直接啃JMS规范文档然后被各种接口绕晕了。其实JMS 1.1的规范文本只有两百多页核心类不超过20个比Spring的代码库友好太多。你先写一个ActiveMQ Demo把生产者、消费者、选择器、事务调一遍再回头翻规范那些文字就不再是抽象的条文而是你刚刚亲眼看到的执行逻辑。最后分享一个小技巧调试JMS相关问题时把Broker的日志级别调到DEBUG然后盯着消息ID在服务端日志里的流转轨迹比单步调试客户端代码更能看清全貌。你能看到消息何时进了队列何时被消费者订阅何时被确认何时被重投。这个视角才是理解JMS标准最直接的入口。