ActiveMQ Classic、JMS 语义与 Broker 故障恢复实战
从一条反复出现的订单消息开始
订单消费者已经把状态写进数据库,几秒后同一条消息却再次到达。开发者在第二次处理时抛出“状态不允许重复变更”,消息持续回滚、重投,最终进入 ActiveMQ.DLQ。有人认为 Broker 重复发送,有人认为事务没有生效,还有人准备直接清空 DLQ。
真正的故障链是:消费者使用事务 session,数据库提交后进程崩溃,JMS session 尚未 commit();Broker 只看到消息没有被确认,于是按至少一次语义重投。重复不是异常现象,而是消息系统与业务数据库之间没有原子事务时必然存在的故障窗口。
排查 ActiveMQ Classic 不能从控制台队列数量开始,也不能把所有 Apache 消息产品统称为 ActiveMQ。应先确认产品线和客户端契约,再沿这条链路定位:
Producer -> JMS Session -> OpenWire Transport -> Classic Broker
-> Destination -> Dispatch/Prefetch -> Consumer Session
-> ACK or Transaction -> KahaDB -> Redelivery -> DLQ先把 Classic 与 Artemis 分清
Apache 同时维护 ActiveMQ Classic 和 ActiveMQ Artemis。两者都能提供 JMS,也都支持多种协议,但并不是同一 Broker 的两个启动模式。
ActiveMQ Classic 延续 5.x 的 Broker、OpenWire、KahaDB、Network of Brokers、Virtual Topic 与 XML 配置模型。当前受支持的发行线是 6.2.x 和 5.19.x;官方 下载页 给出的当前补丁版分别是 6.2.7 与 5.19.8。Classic 6.2 要求 Java 17+,使用 Jakarta Messaging 命名空间;Classic 5.19 要求 Java 11+,主要服务仍依赖 javax.jms 的既有系统。
Artemis 是另一套高性能 Broker 内核,有 address/queue、journal、HA 与集群模型,镜像和配置目录也不同。当前发行版与 Java 要求应从 Artemis 下载页 确认。即使 Artemis 支持 OpenWire,客户端能连通也不等于 Classic 的 broker plugin、KahaDB、Virtual Topic、JMX 属性和 Network of Brokers 配置可以原样迁移。
因此第一项工程判断不是“用 5.x 还是 6.x”,而是:
现有应用导入 javax.jms 还是 jakarta.jms。是否依赖 Classic 专属 destination policy、broker plugin、Virtual Topic、advisory 或 OpenWire 扩展。消息存储是不是 KahaDB,迁移期间是否存在必须保留的 backlog。
团队是在维护既有 Classic 系统,还是为新系统选择长期 Broker 平台。
新项目没有兼容负担时,应把 Artemis、其他现代消息平台与托管服务一起评估;已有 Classic 系统则应先把现状变成可观测、可备份、可回滚的基线,再做迁移。
JMS 的 queue、topic 与 session
JMS 是客户端 API 契约,不是网络协议。Classic Java 客户端通常通过 OpenWire 连接 tcp://host:61616;AMQP、STOMP、MQTT 连接器使用不同协议和端口。把 Web Console 的 8161 写进 brokerURL,或让 AMQP 客户端连 OpenWire 端口,都会出现“端口能通但握手失败”。
Queue:竞争消费和可恢复积压
发送到 queue 的每条消息交给一个消费者。多个消费者用于并行处理,未被确认的消息会在连接中断、session recover 或事务回滚后重新投递。queue 适合作业、命令和需要单个处理结果的任务。
Topic:在线广播与 durable subscription
普通 topic subscriber 只接收在线期间发布的消息。多个 subscriber 各收一份。需要离线恢复时,必须使用 durable subscription,并稳定保存 clientId + subscriptionName;改变其中任意一个都会创建新的订阅状态,旧订阅的 backlog 仍留在 Broker。
Topic 本身不等于事件日志。多个 durable subscriber 会分别持有待消费消息,容量成本随订阅数量和离线时间增长。跨 Broker Network 漂移 durable subscriber 还可能使消息滞留在旧 Broker,因此 Classic 常使用 Virtual Topic 把广播转成可观测的 consumer queue。
Session:顺序、确认与事务边界
JMS Session 是单线程上下文。它承载生产、消费、消息顺序和事务或 ack 状态。不要让多个业务线程并发调用同一个 Session;连接可以复用,Session 通常按工作线程或消费容器管理。
非事务 session 常见确认模式:
AUTO_ACKNOWLEDGE:监听器正常返回或同步 receive 成功返回后,客户端自动确认。处理函数若在外部异步线程继续工作,自动确认可能早于业务完成。CLIENT_ACKNOWLEDGE:应用调用 message.acknowledge()。在 JMS 语义中,它通常确认该 Session 已消费的所有消息,而不是只确认当前一条。DUPS_OK_ACKNOWLEDGE:允许客户端延迟批量确认,提高吞吐但接受更多重复窗口。
事务 session 通过 commit() 原子提交该 Session 中的发送和消费确认,通过 rollback() 撤销并触发重投。官方 事务实现说明 表明 Broker 会缓存事务中的 send 和 ack,直到 commit 才真正生效。JMS 本地事务只覆盖同一个 Session,不自动包含业务数据库。
从零启动一个可观察的 Classic Broker
本地实验选择 apache/activemq:6.2.7,固定补丁版本而不是使用 latest。官方 Docker 入口 映射 61616 用于 OpenWire、8161 用于 Web Console。
先检查 Docker、端口和现有容器:
docker version
Get-NetTCPConnection -LocalPort 61616,8161 -ErrorAction SilentlyContinue
docker ps -a --filter name=activemq-classic-dev启动 Broker:
docker pull apache/activemq:6.2.7
docker run -d --name activemq-classic-dev \
-p 127.0.0.1:61616:61616 \
-p 127.0.0.1:8161:8161 \
-v activemq-classic-data:/opt/apache-activemq/data \
apache/activemq:6.2.7查看启动证据:
docker logs activemq-classic-dev --tail 120
docker exec activemq-classic-dev /opt/apache-activemq/bin/activemq status日志应显示 Broker 启动完成、OpenWire connector 监听 61616。浏览器访问 http://127.0.0.1:8161/admin/ 只能证明 Web 组件可达,不能证明 OpenWire、生产者、消费者和持久化链路正常。
开发镜像中的默认账户只允许本机实验。共享环境必须替换默认口令、拆分应用与管理身份,并限制 Console 网络入口;不能因为绑定了公司内网就继续使用公开默认凭证。
需要验证原生发行包、Windows 服务或 systemd 接管时,从官方 Classic 下载页 取得二进制包并校验 SHA-512/ASC 签名。Linux 的最小启动链如下:
tar -xzf apache-activemq-6.2.7-bin.tar.gz
cd apache-activemq-6.2.7
bin/activemq console另一个终端执行:
bin/activemq status
tail -f data/activemq.log发行包模式把 conf/、data/、lib/ 和 Java 运行时直接交给主机管理,适合既有虚拟机部署和插件兼容验证;容器模式更容易固定镜像与清理环境。生产使用原生包时,应创建独立系统用户、固定 ACTIVEMQ_HOME 与 ACTIVEMQ_BASE、把数据盘与安装目录分离,并由服务管理器发送可控停止信号。不要用管理员账户长期运行 Broker。
Compose 固化端口与数据目录
services:
activemq:
image: apache/activemq:6.2.7
container_name: activemq-classic-dev
ports:
- "127.0.0.1:61616:61616"
- "127.0.0.1:8161:8161"
volumes:
- activemq-data:/opt/apache-activemq/data
healthcheck:
test:
- CMD-SHELL
- /opt/apache-activemq/bin/activemq status | grep -q running
interval: 10s
timeout: 5s
retries: 12
restart: unless-stopped
volumes:
activemq-data:docker compose up -d
docker compose ps
docker compose logs --tail 120 activemq镜像升级前先执行 docker inspect 确认实际挂载路径。路径错位会启动一个空 Broker,看起来“升级成功”,实际旧消息仍留在未挂载或旧 volume 中。
最小正向实验:queue 能发送、接收和归零
Classic 发行包附带命令行 producer/consumer,适合验证 Broker 主链路。先发送一条具有可识别正文的消息:
docker exec activemq-classic-dev \
/opt/apache-activemq/bin/activemq producer \
--destination queue://acme.order.dev.jobs \
--message "eventId=evt-amq-001,orderId=ord-1001" \
--messageCount 1打开 Console 的 Queues 页面,acme.order.dev.jobs 的 enqueue count 应增加,queue size 应为 1。再消费:
docker exec activemq-classic-dev \
/opt/apache-activemq/bin/activemq consumer \
--destination queue://acme.order.dev.jobs \
--messageCount 1预期看到消息正文;Console 中 dequeue count 增加,queue size 回到 0。若 enqueue count 不变,先核对 destination;若 queue size 为 0 但消费者没有打印,检查是否有其他在线 consumer 抢走消息。
这条实验只证明默认客户端完成了一次消费。要验证持久化,必须让消息在没有消费者时跨 Broker 重启仍存在。
对照实验:普通 topic 只广播给在线订阅者
先在终端 A 启动 topic consumer:
docker exec activemq-classic-dev \
/opt/apache-activemq/bin/activemq consumer \
--destination topic://acme.order.dev.events \
--messageCount 1再在终端 B 发布:
docker exec activemq-classic-dev \
/opt/apache-activemq/bin/activemq producer \
--destination topic://acme.order.dev.events \
--message "eventId=evt-topic-online" \
--messageCount 1终端 A 应收到消息。随后反过来操作:在没有任何 subscriber 时先发布 evt-topic-offline,再启动同样的 consumer。新 consumer 不会补收旧消息;等待片刻后按 Ctrl+C 结束。这证明普通 topic subscriber 是在线广播,而不是有位点的事件日志。
需要离线恢复时,客户端要建立 durable subscription,并让身份跨重启稳定:
connection.setClientID("billing-service-prod");
Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE);
var topic = session.createTopic("acme.order.prod.events");
var consumer = session.createDurableSubscriber(topic, "billing-order-events-v1");Broker 用 clientID + subscriptionName 标识这份持久订阅。部署两个副本却复用不兼容的 client ID、随版本随机改变订阅名,或上线新订阅后忘记注销旧订阅,都会制造连接冲突或永久 backlog。删除 durable subscription 前必须确认其 pending、owner 和回放需求,不能只看当前没有在线连接。
正向实验:persistent message 跨重启恢复
JMS 默认 delivery mode 是 persistent。官方 持久与非持久消息说明 指出,persistent 消息写入磁盘或数据库以跨 Broker 重启恢复,non-persistent 在 Broker 故障时可能丢失。
向一个没有消费者的 queue 发送 persistent 消息:
docker exec activemq-classic-dev \
/opt/apache-activemq/bin/activemq producer \
--destination queue://acme.order.dev.restart \
--message "eventId=evt-amq-restart" \
--messageCount 1 \
--persistent true
docker restart activemq-classic-dev
docker logs activemq-classic-dev --tail 120Broker 恢复后,Console 中 queue size 应仍为 1。再消费并确认归零:
docker exec activemq-classic-dev \
/opt/apache-activemq/bin/activemq consumer \
--destination queue://acme.order.dev.restart \
--messageCount 1如果重启后消息消失,依次检查:发送是否确实 persistent、Broker 的 persistent 是否启用、KahaDB 路径是否挂载到 volume、重启是否误用了新容器或新数据目录。
KahaDB:消息为什么还在磁盘里
KahaDB 是 Classic 默认文件持久化存储。官方 KahaDB 文档 将它描述为本地 Broker 使用、面向快速持久化的文件数据库。它不是一条消息一个文件,而是 journal、索引、checkpoint 和清理线程共同维护的追加写存储。
典型配置:
<broker xmlns="http://activemq.apache.org/schema/core"
brokerName="acme-classic"
persistent="true"
deleteAllMessagesOnStartup="false">
<persistenceAdapter>
<kahaDB directory="${activemq.data}/kahadb"
journalMaxFileLength="64mb"
checkpointInterval="5000"
cleanupInterval="30000"
enableIndexWriteAsync="false" />
</persistenceAdapter>
</broker>字段会直接改变故障与性能特征:
directory 决定 Broker 锁、journal 和索引所在位置。两个活动 Broker 不能各自把同一目录当作本地存储。journalMaxFileLength 控制 journal 文件滚动大小。太小增加文件与清理开销,太大增加空间回收粒度。checkpointInterval 控制内存索引状态落盘节奏,不等于消息同步落盘策略。
cleanupInterval 周期性判断哪些 journal 已无未决引用,可以删除。enableIndexWriteAsync 改变索引写入时序,不应在没有恢复演练和性能基线时随意打开。
KahaDB journal 无法删除时,常见原因不是 cleanup 线程坏了,而是文件中仍有一条未确认消息、持久订阅记录、准备中的事务或 DLQ 引用。磁盘增长排查应先看 destination、in-flight、durable subscription、prepared transaction 和日志中的 slow KahaDB access,再判断是否需要离线恢复。不要在 Broker 运行时手工删除 journal 文件。
docker exec activemq-classic-dev sh -lc \
'du -sh /opt/apache-activemq/data && find /opt/apache-activemq/data -maxdepth 3 -type f | sort | head -50'这条命令只用于观察。生产备份必须采用与 Broker 状态一致的方案;直接复制正在写入的目录可能得到不一致快照。
反向实验:事务回滚触发重投
CLI 很难展示完整 Session 事务,下面用 Java/JMS 表达真实逻辑。Classic 6.2 应使用与 Jakarta 命名空间匹配的客户端;Classic 5.19 老应用则保留 javax.jms 客户端。Broker 与客户端版本不必机械完全相同,但协议、JMS API 命名空间和支持矩阵必须在升级前验证。
Classic 6.2 示例的 Maven 依赖为:
<properties>
<maven.compiler.release>17</maven.compiler.release>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.activemq</groupId>
<artifactId>activemq-client</artifactId>
<version>6.2.7</version>
</dependency>
</dependencies>若项目仍导入 javax.jms.*,不要只把依赖版本改成 6.2.7;先留在受支持的 5.19.x 客户端线,或完成 Jakarta 包名和框架兼容迁移。
import jakarta.jms.Connection;
import jakarta.jms.Message;
import jakarta.jms.MessageConsumer;
import jakarta.jms.Queue;
import jakarta.jms.Session;
import org.apache.activemq.ActiveMQConnectionFactory;
public final class RollbackConsumer {
public static void main(String[] args) throws Exception {
var factory = new ActiveMQConnectionFactory("tcp://127.0.0.1:61616");
try (Connection connection = factory.createConnection()) {
connection.start();
Session session = connection.createSession(true, Session.SESSION_TRANSACTED);
Queue queue = session.createQueue("acme.order.dev.rollback");
MessageConsumer consumer = session.createConsumer(queue);
Message message = consumer.receive(5000);
if (message == null) throw new IllegalStateException("no message");
System.out.println("delivery=" + message.getIntProperty("JMSXDeliveryCount"));
session.rollback();
Message redelivered = consumer.receive(5000);
System.out.println("redelivered=" + redelivered.getJMSRedelivered());
System.out.println("delivery=" + redelivered.getIntProperty("JMSXDeliveryCount"));
session.commit();
}
}
}先向 acme.order.dev.rollback 发送一条消息,再运行程序。第一次 delivery count 通常为 1;rollback 后同一消息再次到达,JMSRedelivered=true,delivery count 增加。最终 commit 后 queue size 才归零。
若第二次没有收到,检查消息是否被别的 consumer 抢走、session 是否真的为 transacted、第一次 receive 后是否意外 commit。清理时确认程序已 commit;若实验中断,可在 Console 观察 queue 和 in-flight 状态,再决定消费或 purge。
Ack、事务和业务数据库的故障窗口
一次可靠消费通常包含:读取消息、校验、更新数据库、调用下游、确认消息。JMS 只能控制 Session 内的消息发送和确认,不能自动把普通数据库事务纳入同一原子提交。
常见顺序各有代价:
先 ack/commit JMS,后提交数据库:进程在中间崩溃会永久丢失业务处理机会。先提交数据库,后 ack/commit JMS:进程在中间崩溃会重投,要求数据库操作幂等。XA 事务:把 JMS 与数据库纳入两阶段提交,获得更强原子性,但增加协调、锁持有、故障恢复和运维复杂度。
大多数业务更适合第二种:消息携带稳定 eventId,数据库 inbox 表以 event ID 建唯一键;同一数据库事务写 inbox 与业务状态;事务成功后 commit JMS。重投时唯一键让消费者识别已处理事件并安全确认。
CLIENT_ACKNOWLEDGE 不是“只确认当前消息”的通用替代品。一个 Session 预取多条时,调用某条消息的 acknowledge() 可能确认此前所有已消费消息。希望逐条控制时,应使用事务 session、调整预取,或由成熟的消息监听容器管理边界。
重试、poison message 与 DLQ
Classic 客户端在事务 rollback、未 commit 关闭、Session.recover() 或连接失败后触发重投。官方 重投与 DLQ 文档 说明:超过客户端 redelivery policy 的最大次数后,客户端发送 Poison ACK,Broker 再把消息交给 dead letter strategy。
客户端重投策略示例:
ActiveMQConnectionFactory factory =
new ActiveMQConnectionFactory("tcp://127.0.0.1:61616");
var policy = factory.getRedeliveryPolicy();
policy.setInitialRedeliveryDelay(1000);
policy.setUseExponentialBackOff(true);
policy.setBackOffMultiplier(2.0);
policy.setMaximumRedeliveryDelay(30000);
policy.setMaximumRedeliveries(5);这些参数定义的是技术重试节奏,不应取代业务判断。参数过短会在下游故障时形成重试风暴;参数过长会让消息长期占据 in-flight;最大次数过高会阻塞同一消费者的后续消息并放大日志与外部调用。
默认所有失败消息进入 ActiveMQ.DLQ,多业务共享后难以识别 owner。更可治理的配置是按原 destination 拆分:
<destinationPolicy>
<policyMap>
<policyEntries>
<policyEntry queue="acme.order.>">
<deadLetterStrategy>
<individualDeadLetterStrategy
queuePrefix="DLQ."
useQueueForQueueMessages="true"
processExpired="false" />
</deadLetterStrategy>
</policyEntry>
</policyEntries>
</policyMap>
</destinationPolicy>原队列 acme.order.prod.jobs 的失败消息会进入对应的 DLQ.acme.order.prod.jobs。DLQ consumer 必须记录原 destination、message ID、业务 event ID、redelivery count、异常分类和最后失败原因。修复并重放时要生成审计记录,不能直接在 Console 里批量 move 后不留痕。
反向实验:稳定制造一条 poison message
建立一个 maximumRedeliveries=2 的事务消费者。消费特定 event ID 时始终抛错并 rollback。观察 JMSXDeliveryCount 递增。
达到上限后检查 ActiveMQ.DLQ 或配置的独立 DLQ。
预期证据不是“程序报了三次错”,而是原 queue size/in-flight 变化、DLQ enqueue 增加、消息属性保留原 destination 和失败上下文。实验完成后先导出或记录死信,再通过受控 consumer 消费;不要把 purge 当作清理第一步。
Prefetch、积压和顺序
Classic 客户端会预取消息以提高吞吐。预取后的消息在 Broker 看来处于 dispatched/in-flight,不再计入普通 queue size,却尚未被业务确认。消费者一次预取过多时会出现三个问题:
慢实例拿走大量消息,其他健康实例无事可做。进程故障后大量消息同时重投,形成延迟尖峰。CLIENT_ACKNOWLEDGE 或批量事务的确认边界扩大,重复范围增加。
可以在连接 URI 或 destination policy 中限制 prefetch,例如:
tcp://127.0.0.1:61616?jms.prefetchPolicy.queuePrefetch=50值不是越小越安全。1 有利于严格轮询和长任务,但增加网络往返;较大值适合短小、均匀、幂等的批处理。应依据单条处理 P95、并发 worker、允许的 in-flight 和故障恢复目标测量。
Queue 通常保持 Broker dispatch 顺序,但多个消费者、事务回滚、重投、优先级和 Broker Network 都会改变端到端处理完成顺序。需要同一业务键串行时,可使用 Message Groups(JMSXGroupID)把同组消息黏到一个 consumer,并设置业务 sequence 检查;消费者故障后组会重新分配,仍可能重投。
过期、保留和清理
Classic 不把 queue 当作无限事件日志。消息可以设置 TTL,过期后由 destination policy 决定丢弃或转入 DLQ。KahaDB 只有在 journal 中所有引用都释放后才能回收文件,因此 queue size 已归零不代表磁盘立刻下降。
清理策略要区分:
正常 ack 后的 journal 回收,由 checkpoint/cleanup 管理。消息 TTL 到期,是否进入 DLQ 取决于 dead letter strategy。durable topic subscription 离线积压,必须由订阅 owner 处理或注销。
DLQ 长期积压,需要保留期限、处置责任和重放流程。临时 destination 随连接生命周期销毁,不应用于需要恢复的业务消息。
容量判断至少同时看 enqueue/dequeue 速率、queue size、in-flight、consumer count、oldest message age、DLQ 增长、KahaDB 使用量和磁盘同步延迟。只看 queue size 会漏掉预取消息和 topic durable subscription。
项目接入:把连接、确认和停止顺序写清
Spring Boot 不会因为 YAML 中出现了 spring.activemq.pool.enabled 就自动拥有连接池。项目至少需要 JMS starter;启用池时还必须把 pooled-jms 放进运行时 classpath。下面以 Spring Boot 3、Jakarta Messaging 和 Classic 6.2 为例,依赖版本交给当前 Spring Boot BOM 管理:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-activemq</artifactId>
</dependency>
<dependency>
<groupId>org.messaginghub</groupId>
<artifactId>pooled-jms</artifactId>
</dependency>Classic 5.19 的遗留项目可能仍处于 javax.jms 命名空间,不能只替换 Broker 地址就照搬这组依赖。升级前先检查应用导入、Spring Boot 主版本和 ActiveMQ 客户端变体是否匹配。
配置把 Broker 凭证留在环境变量中,并为开发环境提供明确的 destination:
spring:
activemq:
broker-url: ${ACTIVEMQ_BROKER_URL:tcp://127.0.0.1:61616}
user: ${ACTIVEMQ_USERNAME}
password: ${ACTIVEMQ_PASSWORD}
pool:
enabled: true
max-connections: 4
idle-timeout: 30s
app:
messaging:
order-queue: acme.order.dev.jobs
dlq-prefix: DLQ.
consumer-concurrency: 2-8连接池复用的是 Connection,不应跨线程复用 Session。max-connections 也不是消费者并发数;监听容器会按并发策略创建 Session 和 Consumer。连接上限、监听并发、Broker maximumConnections、prefetch 与下游容量必须一起计算。
先建立一个能看见消息身份的最小收发闭环。生产者不要只发送裸字符串,至少携带稳定业务 ID:
package com.acme.messaging;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.jms.core.JmsTemplate;
import org.springframework.stereotype.Service;
@Service
public class OrderJobProducer {
private final JmsTemplate jmsTemplate;
private final String destination;
public OrderJobProducer(
JmsTemplate jmsTemplate,
@Value("${app.messaging.order-queue}") String destination) {
this.jmsTemplate = jmsTemplate;
this.destination = destination;
}
public void send(String eventId, String body) {
jmsTemplate.convertAndSend(destination, body, message -> {
message.setStringProperty("eventId", eventId);
message.setStringProperty("schemaVersion", "v1");
return message;
});
}
}消费者让异常继续抛给监听容器,不能在 catch 中只记日志后正常返回,否则容器会把失败消息当成成功处理:
package com.acme.messaging;
import jakarta.jms.JMSException;
import jakarta.jms.Message;
import org.springframework.jms.annotation.JmsListener;
import org.springframework.stereotype.Component;
@Component
public class OrderJobListener {
@JmsListener(
destination = "${app.messaging.order-queue}",
concurrency = "${app.messaging.consumer-concurrency}")
public void consume(String body, Message message) throws JMSException {
String eventId = message.getStringProperty("eventId");
int deliveryCount = message.getIntProperty("JMSXDeliveryCount");
System.out.printf("consumed eventId=%s delivery=%d bodyLength=%d%n",
eventId, deliveryCount, body.length());
}
}开发环境可用一个临时 ApplicationRunner 调用 producer.send("evt-boot-001", "create-order")。启动应用后,正向证据应同时成立:日志只出现一次 eventId=evt-boot-001,delivery=1,Console 中 acme.order.dev.jobs 的 enqueue/dequeue 各增加 1,queue size 最终回到 0。只看到发送日志而没有 Broker 计数,不能证明消息真正进入并离开了队列。
反向实验要区分依赖缺失与 Broker 故障:
临时移除 pooled-jms 但保留 pool.enabled=true。应用可能仍能以非池化连接工厂启动,因此要通过 Actuator beans 端点、启动期 bean 类型断言或集成测试确认实际 ConnectionFactory 不是 org.messaginghub.pooled.jms.JmsPoolConnectionFactory;恢复依赖后应重新断言为池化类型。只看“应用启动成功”无法证明连接池生效。
保持应用运行,停止本地 Broker 后再发送 evt-boot-broker-down。发送端应得到 JmsException 或连接异常,不能记录“发送成功”;Broker 恢复后,确认连接池能重新建连,再发送新的 event ID 验证收发。是否自动重试取决于 failover URL 和应用策略,不能把一次方法返回等同于最终业务成功。
让监听器对特定 event ID 抛出异常,检查 JMSXDeliveryCount 是否递增以及消息是否按 redelivery policy 进入独立 DLQ。若日志只报错一次且 queue 直接归零,优先检查异常是否被吞掉、监听容器确认模式和事务边界。
这些实验结束后删除临时 runner,保留生产者、监听器和自动化集成测试。监听容器的事务管理器、ack 模式、receive timeout 和 redelivery policy 要作为同一套配置评审。
优雅停止顺序是:先让实例从流量和健康检查中摘除,暂停接收新消息,等待正在处理的事务完成,commit/rollback 当前 Session,再关闭 consumer、session 和 connection。直接结束 JVM 会让 in-flight 消息重投;如果业务操作非幂等,发布过程本身会制造事故。
生产者必须设置稳定业务 ID、消息类型、schema 版本、correlation ID 和 trace context。不要把密码、token、完整个人信息或大对象放进 message property 和正文;它们可能进入 KahaDB、DLQ、Console、JMX、日志和备份,保留时间远长于一次请求。
认证、授权与传输安全
Classic 可通过 JAAS 或 simple authentication plugin 认证,通过 authorization map 对 queue/topic 的 read、write、admin 授权。官方 安全文档 特别指出,admin 权限允许按需创建 destination,不应默认给所有应用。
<plugins>
<jaasAuthenticationPlugin configuration="activemq-domain" />
<authorizationPlugin>
<map>
<authorizationMap>
<authorizationEntries>
<authorizationEntry queue="acme.order.prod.>"
read="order-consumer"
write="order-producer"
admin="mq-platform" />
<authorizationEntry topic="ActiveMQ.Advisory.>"
read="order-producer,order-consumer"
write="order-producer,order-consumer"
admin="order-producer,order-consumer" />
</authorizationEntries>
</authorizationMap>
</map>
</authorizationPlugin>
</plugins>应用账户、只读排障账户、Broker 管理账户和 Web Console 账户要分离。凭证由密钥系统注入,不写进 Compose、Spring YAML、镜像层、命令历史和截图。轮换前要验证旧新凭证并存窗口以及连接池重连行为。
OpenWire 生产连接应使用 ssl:// 或受保护网络,证书校验不能关闭。Console 经反向代理暴露时,仍需限制来源、强认证、审计 purge/move/delete 等高风险动作。JMX/Jolokia 也属于管理面,不能因为“只有监控读取”就公网暴露。
单机、Broker Network 与高可用不是同一件事
单 Broker + KahaDB
组成最简单:一台 Broker、本地 KahaDB、若干 producer/consumer。它能跨进程重启恢复 persistent 消息,但 Broker 主机或磁盘损坏会中断服务甚至丢数据。适合本地开发、测试和可接受较长恢复时间的内部系统。
Network of Brokers:扩展连接和转发兴趣
Network Connector 在 Broker 间建立 store-and-forward 通道,根据远端消费者兴趣转发消息。官方 Networks of Brokers 文档 说明,连接默认单向,duplex=true 才在一条连接上双向转发。
<networkConnectors>
<networkConnector name="orders-to-hub"
uri="static:(ssl://mq-hub.internal:61617)"
duplex="true"
networkTTL="2"
conduitSubscriptions="true"
decreaseNetworkConsumerPriority="true">
<dynamicallyIncludedDestinations>
<queue physicalName="acme.order.prod.>" />
</dynamicallyIncludedDestinations>
</networkConnector>
</networkConnectors>networkTTL 限制消息跨 Broker 跳数,防止复杂拓扑循环放大。conduitSubscriptions 合并相同订阅兴趣,减少 network consumer 数量,但与 selector 组合时要验证路由。included/excluded destinations 限制跨 Broker 数据边界,避免默认转发所有业务。
decreaseNetworkConsumerPriority 让本地消费者优先,减少不必要跨网转发。
Broker Network 提供扩展和地域转发,不自动复制每个 Broker 的全部 KahaDB。某条消息在任一时刻仍可能只属于一个 Broker;该 Broker 宕机时,消息要等它恢复,除非该 Broker 自身还有 HA 存储形态。
典型故障包括:网络桥断开导致源 Broker 积压;双向拓扑和 selector 使消息滞留错误节点;durable topic subscriber 漂移后旧 Broker 仍持有 backlog;TTL 和 duplicate audit 组合导致消息不再回流。判断时要同时检查每个 Broker 的 queue、network connector 状态、advisory、桥接日志和实际 message ID 流向。
Shared File System Master/Slave:一个逻辑 Broker 的接管
官方 共享文件系统主从说明 描述的机制是:多个 Broker 指向同一 KahaDB 目录,只有取得独占文件锁的实例成为 master 并开放 connector,其余实例等待锁;master 退出后,另一个实例取得锁并接管。
<persistenceAdapter>
<kahaDB directory="/shared/activemq/kahadb" />
</persistenceAdapter>客户端使用 failover transport:
failover:(tcp://mq-a.internal:61616,tcp://mq-b.internal:61616)?randomize=false这种形态的优点是消息状态只有一份,不需要 Broker 间复制;缺点是共享存储成为性能与故障中心,并且文件锁必须在实际文件系统上可靠工作。共享目录可同时挂载不代表 Java 文件锁语义正确。上线前必须做断电、网络分区、锁丢失和接管时间演练;若两个 Broker 同时认为自己持有锁,结果不是高可用,而是数据损坏风险。
JDBC Master/Slave:用数据库独占锁选主
官方 JDBC Master/Slave 文档 描述了另一条共享 store 路径:多个 Broker 共用同一套 JDBC message store。每个实例启动时都尝试取得数据库锁,拿到独占锁的实例成为 master 并开放 connector;其余实例轮询锁,直到 master 退出或锁失效后接管。它解决的是 Broker 进程接管,不会把数据库本身变成高可用。
<persistenceAdapter>
<jdbcPersistenceAdapter
dataSource="#brokerDataSource"
lockKeepAlivePeriod="5000" />
</persistenceAdapter>JDBC 方案借用了团队已有的数据库备份、审计和 HA 能力,也避开了共享文件系统锁语义差异;代价是每次持久消息、确认、事务和调度状态都要经过数据库。高写入吞吐下,数据库日志、索引、锁竞争和网络往返会比本地 KahaDB 更早成为瓶颈。Broker 共用一个单实例数据库时,数据库故障会让所有 Broker 同时失去 store,因此不能把“两台 Broker”误判为端到端高可用。
上线前至少测量持久发送 P95/P99、数据库提交延迟、锁保活失败、master 接管耗时和积压恢复吞吐。还要模拟数据库主从切换:旧 master 是否及时失锁、新 Broker 是否在数据库恢复前错误开放连接、应用 failover 是否发生重复投递。数据库连接池必须为 Broker 锁连接保留容量,不能让业务 store 写满连接池后饿死锁保活。
几种路径的取舍可以落到故障域和团队能力:
| 路径 | 主要优势 | 关键代价与边界 | 更合适的场景 |
|---|---|---|---|
| Shared File System Master/Slave | KahaDB 语义直接、接管后读取同一 store | 依赖低延迟共享存储和可靠文件锁,存储仍是共同故障点 | 已验证 SAN/NFS 锁语义、Classic 存量较重 |
| JDBC Master/Slave | 复用数据库 HA、备份与审计体系 | 持久消息受数据库延迟和锁竞争制约,数据库容量与故障会影响全部 Broker | 吞吐中等、数据库平台成熟、希望集中治理存储 |
| 托管消息服务 | 平台承担节点、存储和升级 | 产品语义、成本、网络与供应商边界需要重新验证 | 团队不希望继续自管 Broker 和共享存储 |
| Artemis | 更现代的 journal、集群和 HA 模型 | 不是 Classic 配置的原位升级,必须迁移语义与 backlog | 新系统或愿意投入协议、客户端和运维迁移 |
Shared FS 与 JDBC 都是共享 store 接管模型,不提供多活写入。若需求是跨地域存活、独立故障域或持续扩展吞吐,应重新评估托管服务、Artemis 或其他消息平台,而不是继续叠加 Classic master/slave。
不要使用已弃用的 Pure Master Slave 或已移除的 LevelDB 方案作为新生产基线。Classic 的 HA 选项应依据当前受支持文档、存储基础设施和恢复演练重新选择。
Broker Network + 每节点 HA
较完整的 Classic 拓扑通常是多个逻辑 Broker 通过 Network Connector 扩展,每个逻辑 Broker 自身由共享存储 master/slave 保护。这样同时承担连接扩展和单 Broker 故障接管,但配置、消息流向、存储、证书和故障演练复杂度显著上升。
当团队只是为了“加两台机器”就引入该拓扑,应比较托管 MQ、Artemis 或其他平台的总成本。Classic 能实现并不代表它仍是新系统最合适的长期选择。
故障诊断:先分类,再动消息
Console 可用,应用连接失败
先区分 Web 与 Broker 端口:
Test-NetConnection 127.0.0.1 -Port 8161
Test-NetConnection 127.0.0.1 -Port 616168161 成功、61616 失败通常是端口未映射、connector 未启动或防火墙问题。两端口都通但握手失败时,检查 URL scheme、客户端协议、TLS 与 Classic/Artemis 产品线。
queue size 不高,但延迟很大
检查 in-flight、consumer prefetch、oldest message age 和业务处理时间。大量消息已预取时,queue size 会显得很低。停止一个消费者观察消息是否重投,可以帮助确认消息到底在 Broker queue 还是客户端缓冲。
消息持续重投
查看 JMSRedelivered、JMSXDeliveryCount、应用异常、事务 commit/rollback 和连接断开。若总在固定处理时长后重投,检查客户端超时、事务超时和下游耗时;若数据库已成功但仍重投,检查 JMS commit 是否位于数据库提交之后、幂等表是否生效。
KahaDB 无法启动
先保留数据目录和完整日志,识别是锁未释放、权限、磁盘已满、journal 缺失、索引损坏还是版本不兼容。不要反复用 deleteAllMessagesOnStartup=true 尝试启动,也不要直接删除 db.data 或 journal。恢复操作应在副本或备份上演练,并记录丢失窗口。
DLQ 快速增长
按 destination、异常类型、producer version、schema version 和首次失败时间聚合。若所有消息同一时间失败,优先排查下游服务、凭证或 schema 发布;若只有少数固定 event ID,通常是数据型 poison message。先隔离生产者或暂停消费风暴,再修复与重放。
备份、恢复、升级与迁移
备份目标必须同时包含 Broker XML、JAAS/授权配置、证书、插件版本、KahaDB 状态与恢复说明。只备份 KahaDB 而缺少 destination policy 和客户端 redelivery 参数,恢复出的消息行为仍可能不同。
获得一份可恢复的 KahaDB 离线备份
最容易犯的错误是在 Broker 持续写入时直接 cp KahaDB 目录。journal、索引、事务和订阅元数据可能处于不同时间点,文件都复制成功也不代表 store 一致。小型 Classic 环境可以使用停写窗口获得离线快照;无法接受停写时,应采用经过验证的存储快照或迁移方案,并以隔离恢复证明其一致性。
先冻结业务入口,停止 producer,再暂停 consumer,让正在处理的 JMS 事务完成。记录一个恢复基线:每个 queue 的 pending/in-flight、每个 durable subscription 的 pending、各 DLQ 数量,以及 prepared transaction。可从 Console/JMX 导出,也可截图并保存 message ID 样本;只记录总消息数无法发现消息落错 destination。
然后优雅停止 Broker,并确认进程已退出:
docker stop --time 60 activemq-classic-dev
docker inspect -f '{{.State.Status}} {{.State.ExitCode}}' activemq-classic-dev只有状态为 exited 后才能复制数据。下面把命名 volume 打成只读归档,并生成校验值;Windows PowerShell 可把 ${PWD} 替换为绝对目录并确保 Docker Desktop 已共享该路径:
mkdir -p backup/activemq-snapshot
docker run --rm \
-v activemq-classic-data:/source:ro \
-v "${PWD}/backup/activemq-snapshot:/backup" \
alpine:3.20 \
sh -c 'cd /source && tar -czf /backup/kahadb.tgz . && sha256sum /backup/kahadb.tgz > /backup/kahadb.sha256'
cp compose.yaml backup/activemq-snapshot/
cp conf/activemq.xml backup/activemq-snapshot/实际目录名应以部署配置为准,证书、JAAS、插件和镜像 digest 也要一并归档;示例路径不存在时不能用空文件替代。归档目录权限应只授予恢复人员,备份中的消息正文、账户和证书都属于敏感数据。
在隔离环境证明备份真的可恢复
恢复不能覆盖原 volume。先校验归档,再写入一个全新的 volume:
cd backup/activemq-snapshot
sha256sum -c kahadb.sha256
cd ../..
docker volume create activemq-classic-restore-test
docker run --rm \
-v activemq-classic-restore-test:/restore \
-v "${PWD}/backup/activemq-snapshot:/backup:ro" \
alpine:3.20 \
sh -c 'cd /restore && tar -xzf /backup/kahadb.tgz'使用与备份时相同的 Classic 镜像和配置启动隔离 Broker,改用本机回环高位端口,禁止连接生产 Broker Network,也不要让生产应用发现它:
docker run -d --name activemq-classic-restore-test \
-p 127.0.0.1:261616:61616 \
-p 127.0.0.1:28161:8161 \
-v activemq-classic-restore-test:/opt/apache-activemq/data \
apache/activemq:6.2.7
docker logs activemq-classic-restore-test --tail 200启动成功只是第一关。按停机前基线逐项核对:
queue 名称、pending 数量和抽样 message ID 是否与基线一致,停机前 in-flight 是否已经归零。durable subscription 的 clientId + subscriptionName、active/inactive 状态和 pending 是否一致;不要启动真实 subscriber 抢走恢复数据。ActiveMQ.DLQ 及独立 DLQ.* 的数量、原 destination 和 event ID 是否一致。
prepared transaction 是否存在,决定恢复、提交或回滚前先确认业务系统状态。用隔离测试账户消费一条专门预留的验证消息,确认 KahaDB 不只是能打开,还能完成 dispatch 与 ack;随后销毁隔离 Broker 和测试 volume。
任何计数差异、journal 缺失、索引重建失败或 durable subscription 消失,都意味着这份备份不能进入恢复资产库。先保留日志和归档调查,不能把“Broker 最终启动了”作为通过标准。
docker stop activemq-classic-restore-test
docker rm activemq-classic-restore-test
docker volume rm activemq-classic-restore-test生产恢复时仍应创建新数据目录并先隔离验收,再切换连接入口;保留原 store 为只读证据和回退点。恢复期间禁止 producer 同时向新旧 Broker 写入,切换后用稳定 event ID 和 inbox 幂等处理可解释的重复投递。
恢复演练至少验证:
未消费 persistent message 在 Broker 重启后存在。master 失效后客户端通过 failover URL 重连,未确认消息发生可解释的重投。KahaDB 从一致性备份恢复后,queue、durable subscription、DLQ 与 prepared transaction 数量可核对。
回滚 Broker 版本时,旧版本是否能读取已升级的 store;不能假设存储格式永远双向兼容。
从 Classic 迁移 Artemis 时,先做功能清单而不是只测 OpenWire 连接:
queue/topic、durable subscription、selector、message group 和事务语义是否一致。Virtual Topic、advisory、broker plugin、Network Connector 和 JMX 依赖如何替换。KahaDB backlog 是停机导出导入,还是双 Broker 在线搬迁。
producer/consumer 切流期间如何防重复、如何回退、旧 Broker 何时只读。Classic 5.x 的 javax.jms 客户端是否升级,还是先通过协议兼容过渡。
Apache 的 Classic/Artemis 迁移说明 明确提示 Artemis 并非 Classic 的百分之百重实现,并提供 KahaDB 导出与线上迁移方向。迁移验收必须基于业务语义和故障演练,不能以“客户端成功连接 Artemis”结束。
清理和回滚
本地实验先停止生产者和消费者,再停止 Broker:
docker stop activemq-classic-dev
docker rm activemq-classic-dev保留 volume 可以用同版本容器重新挂载恢复。确认所有实验消息都可删除后,才执行:
docker volume rm activemq-classic-dataCompose 环境使用:
docker compose down # 保留数据
docker compose down -v # 删除数据,仅用于确认无保留价值的本地环境共享环境不允许通过删除 volume、KahaDB 目录或 Console purge 完成“回滚”。配置变更要保留旧镜像、旧 XML、凭证兼容窗口与客户端 failover 路径;消息变更要有 move/export/replay 记录;迁移失败时要能停止新 Broker 写入并把 producer/consumer 切回旧 Broker,而不会让两边同时消费产生双写。
容量、成本和长期治理
Classic 的成本不只由消息吞吐决定。Persistent send 的同步磁盘确认、事务批次、KahaDB 清理、durable topic backlog、DLQ 保留、Broker Network 跨网流量、共享存储延迟和 Console/JMX 安全都会进入总成本。
容量预算至少包含:
未消费容量 = 峰值写入速率 × 平均消息大小 × 最大恢复时间
持久订阅容量 = 每个 durable subscriber 的离线窗口之和
磁盘预算 = 未消费 + DLQ + journal 回收滞后 + 事务/索引开销 + 安全余量
恢复吞吐 = 可用消费吞吐 - 实时新增吞吐如果恢复吞吐不大于实时新增吞吐,积压永远追不平。阈值应由业务 SLO、基线压测和磁盘增长趋势确定,而不是复制一个固定 queue depth。
团队职责需要落到对象:
平台 owner 负责 Broker、KahaDB、connector、Network、HA、证书、备份恢复和升级。destination owner 负责 queue/topic 命名、schema、TTL、重试、DLQ 和容量预算。应用 owner 负责事务边界、幂等、prefetch、并发、优雅停止和重放安全。
安全 owner 负责账户、最小权限、Console/JMX 暴露、凭证轮换和敏感数据保留。
上线评审应能给出证据:persistent 消息跨重启恢复;rollback 会重投且业务幂等;poison message 可进入独立 DLQ;单 Broker 故障后的客户端行为已演练;KahaDB 磁盘触顶前能报警;管理动作可审计;Classic 与 Artemis 的版本、API 命名空间和迁移边界已经明确。
当系统能回答“消息此刻在哪个 Broker、哪个 destination、哪个 Session 或哪个 KahaDB 引用中,失败后由谁以什么证据恢复”,ActiveMQ Classic 才从一个遗留中间件变成了可维护的工程基础设施。
