JMS、ActiveMQ 学习一则:从连接工厂到消息收发的完整链路拆解

发布时间:2026/10/8 6:16:17
JMS、ActiveMQ 学习一则:从连接工厂到消息收发的完整链路拆解 1. 从一次本地消息“发出去没人收”说起JMS 与 ActiveMQ 的完整链路到底长什么样如果你刚开始接触 JMS 和 ActiveMQ大概率会遇到这样一个场景代码里producer.send()明明没报错日志也打印了“发送成功”但消费者那边就是一动不动队列里也看不到消息。我第一次搭本地环境时就卡在这里后来才发现问题出在ConnectionFactory的连接参数和Destination的名字对不上——一个发到了queue.test另一个却在监听test.queue。JMSJava Message Service是一套 Java 消息中间件的规范接口它定义了ConnectionFactory、Destination、Producer、Consumer这几个核心对象的标准用法但不关心底层是谁实现的。ActiveMQ 就是这套规范的一个经典实现你可以把它理解成“JMS 规范的一份可运行答案”。规范负责约定“怎么连、怎么发、怎么收”ActiveMQ 负责“真正把消息存下来、转发出去”。这篇文章面向的是想独立跑通一条消息链路的开发者。我会围绕四个核心对象把点对点Queue和发布订阅Topic两种模型的消息流转路径拆开讲给出可以直接复制的 Java 示例和本地启动验证步骤。你跟着做完应该能亲眼看到一条消息从生产者发出、进入 Broker、再被消费者取走的完整过程而不是停留在“概念懂了但跑不起来”的状态。需要说明的是本文的示例全部基于本地 ActiveMQ Broker不涉及任何网络穿透或非正规访问方式。如果你后续想把消息能力接到自己的 AI 应用或编码工作流里也可以参考 TaoToken 这类平台提供的模型接入能力把消息触发和模型调用串起来这部分我会在最后一节给出衔接思路。2. 前置准备ActiveMQ 本地启动与 TaoToken 接入前的环境确认在写第一行 Java 代码之前先把“地基”打好。这一节的目标是让你拥有一个正在运行的 ActiveMQ Broker并且确认 Java 和 Maven 环境可用。很多人跑不通链路不是代码问题而是 Broker 根本没起来或者端口被占用。2.1 下载与启动 ActiveMQ BrokerActiveMQ 的经典版本是 5.x 系列对 JMS 1.1 规范支持完整适合入门。下载解压后目录结构大致如下apache-activemq-5.18.3/ ├── bin/ │ ├── activemq # Linux/Mac 启动脚本 │ └── activemq.bat # Windows 启动脚本 ├── conf/ │ └── activemq.xml # 核心配置定义 Broker、传输协议、持久化 └── lib/启动命令Linux/Maccd apache-activemq-5.18.3/bin ./activemq startWindows 下用cd apache-activemq-5.18.3\bin activemq.bat start启动成功后默认会监听几个关键端口61616是 OpenWire 协议端口Java 客户端默认连这个8161是 Web 控制台端口。打开浏览器访问http://localhost:8161/admin默认账号密码都是admin。如果你能看到 Queues 和 Topics 的管理页面说明 Broker 已经就绪。注意如果 8161 打不开先检查conf/jetty.xml里的端口配置以及本机是否有其他程序占用了 8161。启动日志在data/activemq.log报错信息都在里面。2.2 确认 Java 与 Maven 环境ActiveMQ 5.x 需要 JDK 8 及以上。执行java -version mvn -version如果 Maven 没装可以用 IDE 自带的或者手动配置。接下来创建一个 Maven 项目在pom.xml里加入 ActiveMQ 客户端依赖dependency groupIdorg.apache.activemq/groupId artifactIdactivemq-client/artifactId version5.18.3/version /dependency这个依赖里包含了 JMS APIjavax.jms.*和 ActiveMQ 的具体实现。注意JMS 1.1 用的是javax.jms包而 JMS 2.0 用的是jakarta.jms两者不能混用。本文统一用javax.jms和 ActiveMQ 5.x 默认行为一致。2.3 关于 TaoToken 的定位说明如果你后续想把消息消费和 AI 模型调用结合起来比如消费者收到消息后触发一次模型推理那么你需要一个稳定的模型接入入口。TaoToken 提供的是模型 API 接入能力官网是https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_contentAPI 地址是https://taotoken.net/api。它的作用不是替代 ActiveMQ而是当你的消息链路需要“消费后调用模型”时提供一个统一的 Key 和 Base URL 来发起请求。这一点在最后一节会展开。现在Broker 在跑依赖已加可以进入核心对象的拆解了。3. 四个核心对象怎么配ConnectionFactory、Destination、Producer、Consumer 的可复制配置JMS 编程模型里这四个对象是有明确分工的。ConnectionFactory负责“怎么连到 Broker”Destination负责“消息发到哪”Producer负责“把消息放进去”Consumer负责“把消息取出来”。理解它们的关系比死记 API 更重要。3.1 ConnectionFactory连接参数的集中管理ConnectionFactory是客户端和 Broker 之间的入口。ActiveMQ 提供了ActiveMQConnectionFactory最常用的构造参数就是 Broker 的 URLimport org.apache.activemq.ActiveMQConnectionFactory; import javax.jms.ConnectionFactory; ConnectionFactory factory new ActiveMQConnectionFactory( tcp://localhost:61616 );这里的tcp://localhost:61616就是 OpenWire 协议的地址。如果你在activemq.xml里改了端口这里要同步改。生产环境通常还会设置用户名密码ActiveMQConnectionFactory factory new ActiveMQConnectionFactory(admin, admin, tcp://localhost:61616);一个容易被忽略的点是ConnectionFactory本身是线程安全的可以在多个线程间共享但由它创建的Connection不是。所以常见做法是全局维护一个ConnectionFactory实例每次需要时再创建Connection。3.2 DestinationQueue 与 Topic 的分叉点Destination是一个接口它有两个主要子类型Queue和Topic。这两个类型直接决定了消息是“点对点”还是“发布订阅”。在 JMS 1.1 里创建 Destination 有两种方式。一种是通过Session创建Session session connection.createSession(false, Session.AUTO_ACKNOWLEDGE); Queue queue session.createQueue(queue.order); Topic topic session.createTopic(topic.notice);另一种是直接用ActiveMQQueue/ActiveMQTopic构造import org.apache.activemq.command.ActiveMQQueue; import org.apache.activemq.command.ActiveMQTopic; Destination queue new ActiveMQQueue(queue.order); Destination topic new ActiveMQTopic(topic.notice);两种方式效果一样区别在于前者依赖Session后者更独立。名字必须完全一致大小写敏感。我踩过的坑就是生产者用queue.order消费者用queue.Order结果消息进了队列但消费者永远收不到。3.3 Producer 与 Consumer发送和接收的对称结构MessageProducer通过session.createProducer(destination)创建然后调用send()MessageProducer producer session.createProducer(queue); TextMessage message session.createTextMessage(hello jms); producer.send(message);MessageConsumer通过session.createConsumer(destination)创建然后调用receive()或设置MessageListenerMessageConsumer consumer session.createConsumer(queue); Message msg consumer.receive(5000); if (msg instanceof TextMessage) { System.out.println(((TextMessage) msg).getText()); }receive(5000)表示最多等 5 秒超时返回 null。如果不带参数receive()会一直阻塞。生产环境更推荐用监听器方式避免线程被长时间挂起。3.4 一份完整的可复制配置片段把上面的内容整合成一个可直接运行的 Java 类包含发送和接收两个方法import org.apache.activemq.ActiveMQConnectionFactory; import javax.jms.*; public class JmsDemo { private static final String BROKER_URL tcp://localhost:61616; private static final String QUEUE_NAME queue.order; public static void main(String[] args) throws Exception { sendMessage(); receiveMessage(); } static void sendMessage() throws Exception { ConnectionFactory factory new ActiveMQConnectionFactory(BROKER_URL); try (Connection connection factory.createConnection()) { connection.start(); Session session connection.createSession(false, Session.AUTO_ACKNOWLEDGE); Destination destination session.createQueue(QUEUE_NAME); MessageProducer producer session.createProducer(destination); TextMessage message session.createTextMessage(order-1001-created); producer.send(message); System.out.println(sent: message.getText()); } } static void receiveMessage() throws Exception { ConnectionFactory factory new ActiveMQConnectionFactory(BROKER_URL); try (Connection connection factory.createConnection()) { connection.start(); Session session connection.createSession(false, Session.AUTO_ACKNOWLEDGE); Destination destination session.createQueue(QUEUE_NAME); MessageConsumer consumer session.createConsumer(destination); Message msg consumer.receive(5000); if (msg instanceof TextMessage) { System.out.println(received: ((TextMessage) msg).getText()); } else { System.out.println(no message within timeout); } } } }这段代码里connection.start()是关键。JMS 规范规定Connection创建后处于“停止”状态必须调用start()才会真正开始投递消息。很多人忘了这一步结果消费者一直收不到。如果你希望把这段消息链路和模型调用结合比如消费者收到order-1001-created后去请求一次模型那么可以在receiveMessage()里追加 HTTP 调用Base URL 用https://taotoken.net/apiKey 从控制台获取。这样消息系统和 AI 能力就串起来了。4. 跑通验证点对点与发布订阅两种模型的实际请求与结果配置写完了接下来要亲眼看到消息流转。这一节分别验证 Queue 和 Topic 两种模型并给出预期输出。只有看到控制台打印出对应结果才算真正跑通。4.1 点对点模型一条消息只被一个消费者取走点对点模型的核心特征是消息进入 Queue 后只会被一个消费者消费。即使你启动多个消费者同一条消息也只会落到其中一个手里。这适合任务分发场景比如订单处理。验证步骤第一步先启动一个消费者让它阻塞等待。你可以把上面的receiveMessage()单独跑起来或者写一个带MessageListener的版本MessageConsumer consumer session.createConsumer(destination); consumer.setMessageListener(msg - { try { if (msg instanceof TextMessage) { System.out.println(consumer-A got: ((TextMessage) msg).getText()); } } catch (JMSException e) { e.printStackTrace(); } }); Thread.sleep(30000);第二步运行发送端发送三条消息for (int i 1; i 3; i) { TextMessage message session.createTextMessage(order- i); producer.send(message); }预期结果消费者 A 依次打印order-1、order-2、order-3。如果你再启动一个消费者 B重新发送三条消息那么 A 和 B 会各自分到一部分但同一条消息不会同时被两者收到。在 Web 控制台的 Queues 页面你能看到queue.order的Enqueue和Dequeue计数。发送三条后 Enqueue 为 3消费三条后 Dequeue 为 3Pending归零。这个数字变化就是链路跑通的最直接证据。4.2 发布订阅模型一条消息被所有订阅者收到Topic 模型和 Queue 正好相反一条消息会被所有当前活跃的订阅者收到。注意“当前活跃”这个限定JMS 1.1 的 Topic 默认不持久化订阅者不在线就收不到。验证步骤第一步启动两个订阅者都监听同一个 TopicDestination topic session.createTopic(topic.notice); MessageConsumer consumer session.createConsumer(topic); consumer.setMessageListener(msg - { if (msg instanceof TextMessage) { System.out.println(subscriber got: ((TextMessage) msg).getText()); } });第二步等两个订阅者都进入监听状态后发送一条消息Destination topic session.createTopic(topic.notice); MessageProducer producer session.createProducer(topic); producer.send(session.createTextMessage(system-maintenance));预期结果两个订阅者的控制台都会打印system-maintenance。如果你在发送之后才启动第三个订阅者它是收不到这条消息的因为消息已经投递完毕。4.3 两种模型的对照表维度Queue点对点Topic发布订阅消息份数一条消息一份一条消息每个订阅者一份消费者数量多个消费者竞争每个订阅者都收到离线消息默认保留上线后消费默认不保留典型场景订单处理、任务分发通知广播、事件推送Destination 创建session.createQueue()session.createTopic()这张表建议收藏实际选型时先问自己“这条消息是需要被处理一次还是需要被多方感知”答案基本就出来了。4.4 用 TaoToken 模型对话做一次链路延伸验证如果你想让验证更有意思一点可以在消费者收到消息后调用一次模型对话接口把消息内容作为输入。TaoToken 的模型对话入口是https://taotoken.net/api你需要在控制台创建一个 API Key然后在消费者里发起 HTTP 请求。这样你就能看到“消息触发 → 模型响应”的完整链路而不只是控制台打印。这一步不是必须的但它能帮你理解消息系统在 AI 工作流里的位置消息负责解耦和触发模型负责处理内容。5. 本篇常见错误排查401、local proxy failed、reading choices 与 OAuth 报错对照链路跑不通时报错信息往往很具体。这一节把入门阶段最常见的几类错误列出来对照排查。需要强调的是下面涉及的连接问题都应在本地或合规网络环境下解决不涉及任何非正规访问手段。5.1 401 Unauthorized认证失败如果你在连接 Broker 或调用模型接口时看到 401通常意味着凭证不对。ActiveMQ 默认账号密码是admin/admin但如果你改过conf/users.properties和conf/credentials.properties就要用改后的。模型接口的 401 则通常是 API Key 缺失或写错检查请求头里的Authorization字段是否带了正确的 Key。排查顺序先确认用户名密码再确认 Key 是否有多余空格最后确认请求头格式。5.2 local proxy failed本地连接被拦截这个报错一般出现在客户端尝试连接 Broker 或外部接口时本机网络配置导致连接被拦截。解决思路是检查本机的网络设置确认没有异常的本地转发规则。对于 ActiveMQ先确认tcp://localhost:61616是否真的在监听netstat -an | grep 61616如果没有输出说明 Broker 没起来或端口不对。对于模型接口确认https://taotoken.net/api能正常访问可以用curl做一次连通性测试。5.3 reading choices响应解析失败这个报错常见于调用模型接口时客户端期望的响应格式和实际返回不一致。比如你用的是 OpenAI 兼容格式但返回体里没有choices字段。排查方法是先把原始响应打印出来看看到底返回了什么。常见原因包括请求体里的model字段写错、messages结构不对、或者 Base URL 拼错了路径。一个实用的调试片段HttpResponseString response client.send(request, HttpResponse.BodyHandlers.ofString()); System.out.println(status: response.statusCode()); System.out.println(body: response.body());先看 status再看 body比盲目改代码有效得多。5.4 OAuth 相关报错令牌获取失败如果你在接入某些需要 OAuth 的服务时看到令牌相关报错先确认令牌是否过期、回调地址是否配置正确。对于 ActiveMQ 本身它默认不走 OAuth所以这类报错通常出现在外围的模型接入或工具链里。处理原则是先确认令牌有效期再确认请求的 scope 是否包含所需权限。5.5 三件套检查清单无论你用的是 CC Switch、Cline MCP 还是 Codex 的auth.json只要涉及模型接入都要确认三件套齐全Base URLhttps://taotoken.net/apiAPI Key从控制台创建注意保密Model ID按文档填写不要臆造这三项缺任何一项都会导致请求失败。把这三项写进配置文件时注意 JSON 或 TOML 的格式正确性少一个逗号都会解析失败。6. 从消息链路到长期编码工作流把验证过的配置沉淀下来跑通一条消息链路只是起点。真正有价值的是把这套配置沉淀成可复用的模板让下次搭环境时不用从头再来。我的做法是把ConnectionFactory的 URL、账号密码抽到配置文件里把 Queue 和 Topic 的名字定义成常量把发送和接收封装成工具方法。这样业务代码里只需要调用send(queueName, content)和onMessage(queueName, handler)不用每次都写一遍 JMS 样板代码。如果你打算把消息消费和模型调用长期结合起来比如做一个“消息触发模型处理”的 Agent 工作流那么可以考虑使用 Coding Plan 这类长期方案把模型调用额度、Key 管理和编码工作流统一起来。它的入口在 TaoToken 的控制台里可以找到适合需要持续调用模型的场景而不是一次性验证。最后给一个实用建议每次改完activemq.xml或客户端配置先用 Web 控制台确认 Broker 状态再跑最小发送接收用例。不要一上来就跑复杂业务逻辑否则报错时你分不清是配置问题还是业务代码问题。把最小链路跑通再往上叠功能这是我在消息中间件上踩过坑之后最想分享的一条经验。