简介
ActiveMQ 是 Apache 提供的消息中间件。
它在 Java 生态里很常见,尤其是传统企业系统、老项目改造、JMS 规范相关系统中出现得比较多。
简单理解:
1生产者 2 | 3 v 4ActiveMQ Broker 5 | 6 v 7消费者 8
ActiveMQ 最常见的使用方式是 JMS。
JMS 全称是 Java Message Service,它是一套 Java 消息服务 API 规范。
ActiveMQ 和 JMS 的关系可以类比成:
1JMS 是规范 2ActiveMQ 是实现 3
一句话概括:
1ActiveMQ 是 Java 生态里经典的 JMS 消息中间件,适合异步处理、系统解耦、任务队列和发布订阅场景。 2
ActiveMQ 解决什么问题
一个业务动作经常会带出很多后续动作。
比如订单创建成功后:
1扣库存 2发短信 3发邮件 4写日志 5生成物流任务 6同步积分 7
如果全部同步执行,接口会变慢。
如果短信服务异常,订单接口也可能被影响。
使用 ActiveMQ 后,可以拆成:
1订单服务 2 | 3 v 4发送订单消息 5 | 6 v 7ActiveMQ 8 | 9 +--> 库存服务消费消息 10 +--> 短信服务消费消息 11 +--> 邮件服务消费消息 12 +--> 日志服务消费消息 13
常见作用:
- 异步处理
- 系统解耦
- 削峰填谷
- 任务分发
- 发布订阅
- 消息重试和失败兜底
ActiveMQ Classic 和 Artemis
ActiveMQ 需要先区分两条线:
| 名称 | 说明 |
|---|---|
| ActiveMQ Classic | 传统 ActiveMQ 5.x / 6.x 线,老项目和 JMS 系统常见 |
| ActiveMQ Artemis | 新一代消息 Broker,性能和架构更现代 |
简单理解:
1ActiveMQ Classic:经典 JMS Broker 2ActiveMQ Artemis:新一代多协议、高性能 Broker 3
如果是维护老系统,经常会遇到 ActiveMQ Classic。
如果是新系统选型,并且明确要用 Apache ActiveMQ 体系,Artemis 更值得优先评估。
本文示例以 ActiveMQ Classic + Spring Boot JMS 为主,因为它更贴近日常 Java 老系统和企业项目集成场景。
ActiveMQ 和 RabbitMQ 的区别
| 对比项 | ActiveMQ | RabbitMQ |
|---|---|---|
| 常见协议 | JMS、OpenWire、STOMP、AMQP、MQTT | AMQP 为主,也支持多协议 |
| Java 规范关系 | JMS 体系很强 | 不以 JMS 为核心 |
| 核心模型 | Queue / Topic | Exchange / Queue / Binding |
| 路由能力 | JMS 模型更直接 | Exchange 路由更灵活 |
| 常见场景 | 传统企业系统、JMS 项目 | 微服务消息、灵活路由、事件分发 |
ActiveMQ 更像:
1Queue / Topic 模型清晰,适合 JMS 项目。 2
RabbitMQ 更像:
1Exchange + Routing Key 更灵活,适合复杂路由。 2
JMS 核心概念
| 概念 | 说明 |
|---|---|
| ConnectionFactory | 连接工厂 |
| Connection | 客户端到 Broker 的连接 |
| Session | 会话,负责生产、消费、确认、事务 |
| Destination | 目的地,Queue 或 Topic |
| Queue | 点对点队列 |
| Topic | 发布订阅主题 |
| MessageProducer | 消息生产者 |
| MessageConsumer | 消息消费者 |
| Message | 消息对象 |
整体流程:
1ConnectionFactory 2 | 3 v 4Connection 5 | 6 v 7Session 8 | 9 +--> MessageProducer -> Destination 10 | 11 +--> MessageConsumer <- Destination 12
Queue 点对点模式
Queue 是点对点模式。
特点:
1一条消息只会被一个消费者处理。 2
示例:
1order.queue 2 | 3 +--> consumer-1 4 +--> consumer-2 5
如果 order.queue 有 10 条消息,两个消费者会竞争消费。
每条消息最终只会被其中一个消费者处理。
适合:
- 订单处理
- 短信发送
- 邮件发送
- 报表生成
- 后台任务
Topic 发布订阅模式
Topic 是发布订阅模式。
特点:
1一条消息可以被多个订阅者收到。 2
示例:
1order.topic 2 | 3 +--> sms-consumer 4 +--> email-consumer 5 +--> log-consumer 6
订单支付成功消息发送到 order.topic 后,短信、邮件、日志服务都能收到。
适合:
- 业务事件广播
- 通知消息
- 日志订阅
- 缓存刷新
Topic 持久订阅和非持久订阅
Topic 需要注意订阅方式。
| 类型 | 离线期间消息是否保留 | 说明 |
|---|---|---|
| 非持久订阅 | 不保留 | 消费者在线时才能收到 |
| 持久订阅 | 保留 | 消费者离线后再上线也能收到 |
非持久订阅适合临时通知。
持久订阅适合关键业务事件。
持久订阅通常需要设置唯一的 clientId 和订阅名。
Docker 启动 ActiveMQ Classic
本地学习可以使用 Docker。
不同镜像维护方式不同,实际项目以团队镜像仓库或官方发布包为准。
示例:
1docker run -d \ 2 --name activemq-classic \ 3 -p 61616:61616 \ 4 -p 8161:8161 \ 5 apache/activemq-classic:latest 6
端口说明:
| 端口 | 作用 |
|---|---|
| 61616 | JMS / OpenWire 连接端口 |
| 8161 | Web 管理控制台 |
管理后台:
1http://localhost:8161/admin 2
常见默认账号:
1admin / admin 2
如果镜像没有提供默认账号,需要查看镜像说明或自定义配置。
原生 JMS 依赖
Spring Boot 3 使用 jakarta.jms 包名。
如果使用 ActiveMQ Classic 6 或 Jakarta 体系,可以选择 Jakarta 版本客户端。
1<dependency> 2 <groupId>org.apache.activemq</groupId> 3 <artifactId>activemq-client-jakarta</artifactId> 4 <version>${activemq.version}</version> 5</dependency> 6
如果是老项目,可能仍然使用 javax.jms 和 ActiveMQ 5.x 客户端。
1<dependency> 2 <groupId>org.apache.activemq</groupId> 3 <artifactId>activemq-client</artifactId> 4 <version>${activemq.version}</version> 5</dependency> 6
包名差异:
1老项目:javax.jms.* 2新项目:jakarta.jms.* 3
原生 JMS:发送 Queue 消息
1package com.example.activemq.raw; 2 3import jakarta.jms.Connection; 4import jakarta.jms.ConnectionFactory; 5import jakarta.jms.DeliveryMode; 6import jakarta.jms.Destination; 7import jakarta.jms.MessageProducer; 8import jakarta.jms.Session; 9import jakarta.jms.TextMessage; 10import org.apache.activemq.ActiveMQConnectionFactory; 11 12public class RawQueueProducer { 13 14 private static final String BROKER_URL = "tcp://localhost:61616"; 15 private static final String QUEUE_NAME = "order.queue"; 16 17 public static void main(String[] args) throws Exception { 18 ConnectionFactory factory = new ActiveMQConnectionFactory(BROKER_URL); 19 20 try (Connection connection = factory.createConnection(); 21 Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE)) { 22 23 Destination destination = session.createQueue(QUEUE_NAME); 24 MessageProducer producer = session.createProducer(destination); 25 producer.setDeliveryMode(DeliveryMode.PERSISTENT); 26 27 TextMessage message = session.createTextMessage("订单创建成功"); 28 message.setStringProperty("bizType", "order"); 29 message.setStringProperty("messageId", "msg-1001"); 30 31 producer.send(message); 32 33 System.out.println("消息已发送:" + message.getText()); 34 } 35 } 36} 37
说明:
1Session.AUTO_ACKNOWLEDGE:自动确认 2DeliveryMode.PERSISTENT:持久化消息 3createQueue:创建 Queue 目的地 4
原生 JMS:消费 Queue 消息
1package com.example.activemq.raw; 2 3import jakarta.jms.Connection; 4import jakarta.jms.ConnectionFactory; 5import jakarta.jms.MessageConsumer; 6import jakarta.jms.Session; 7import jakarta.jms.TextMessage; 8import org.apache.activemq.ActiveMQConnectionFactory; 9 10public class RawQueueConsumer { 11 12 private static final String BROKER_URL = "tcp://localhost:61616"; 13 private static final String QUEUE_NAME = "order.queue"; 14 15 public static void main(String[] args) throws Exception { 16 ConnectionFactory factory = new ActiveMQConnectionFactory(BROKER_URL); 17 18 Connection connection = factory.createConnection(); 19 Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE); 20 MessageConsumer consumer = session.createConsumer(session.createQueue(QUEUE_NAME)); 21 22 connection.start(); 23 24 consumer.setMessageListener(message -> { 25 try { 26 if (message instanceof TextMessage textMessage) { 27 System.out.println("收到消息:" + textMessage.getText()); 28 } 29 30 message.acknowledge(); 31 } catch (Exception e) { 32 throw new RuntimeException(e); 33 } 34 }); 35 36 System.out.println("消费者已启动"); 37 Thread.currentThread().join(); 38 } 39} 40
CLIENT_ACKNOWLEDGE 表示客户端手动确认。
注意:
1JMS 的 acknowledge() 可能确认当前 Session 中已消费的一批消息。 2
如果需要更强的边界控制,可以考虑事务会话。
JMS 确认模式
JMS 常见确认模式:
| 模式 | 说明 |
|---|---|
| AUTO_ACKNOWLEDGE | 监听器正常返回后自动确认 |
| CLIENT_ACKNOWLEDGE | 客户端调用 acknowledge() 确认 |
| DUPS_OK_ACKNOWLEDGE | 允许延迟确认,可能出现重复 |
| SESSION_TRANSACTED | 使用事务提交或回滚 |
AUTO_ACKNOWLEDGE 使用简单。
关键业务更常见的是事务或手动确认。
JMS 事务
创建事务会话:
1Session session = connection.createSession(true, Session.SESSION_TRANSACTED); 2
处理成功:
1session.commit(); 2
处理失败:
1session.rollback(); 2
示例:
1MessageConsumer consumer = session.createConsumer(session.createQueue("order.queue")); 2 3TextMessage message = (TextMessage) consumer.receive(); 4 5try { 6 System.out.println("处理消息:" + message.getText()); 7 8 // 执行业务逻辑 9 10 session.commit(); 11} catch (Exception e) { 12 session.rollback(); 13} 14
事务回滚后,消息会重新投递。
Spring Boot 依赖
Spring Boot 集成 ActiveMQ Classic:
1<dependency> 2 <groupId>org.springframework.boot</groupId> 3 <artifactId>spring-boot-starter-activemq</artifactId> 4</dependency> 5
如果选择 Artemis:
1<dependency> 2 <groupId>org.springframework.boot</groupId> 3 <artifactId>spring-boot-starter-artemis</artifactId> 4</dependency> 5
Classic 和 Artemis 不建议在同一个应用里混用。
application.yml
ActiveMQ Classic 示例:
1spring: 2 activemq: 3 broker-url: tcp://localhost:61616 4 user: admin 5 password: admin 6 jms: 7 listener: 8 acknowledge-mode: auto 9
如果需要连接池,可以配置连接池依赖和对应参数。
不同 Spring Boot 版本对连接池自动配置细节略有差异,实际项目以当前版本文档为准。
开启 JMS 注解
使用 @JmsListener 需要开启 JMS 注解支持。
1package com.example.activemq.config; 2 3import org.springframework.context.annotation.Configuration; 4import org.springframework.jms.annotation.EnableJms; 5 6@Configuration 7@EnableJms 8public class JmsConfig { 9} 10
消息对象
1package com.example.activemq.order; 2 3import java.math.BigDecimal; 4import java.time.LocalDateTime; 5 6public record OrderCreatedMessage( 7 Long orderId, 8 String orderNo, 9 Long userId, 10 BigDecimal amount, 11 LocalDateTime createdAt 12) { 13} 14
JSON 消息转换器
业务系统里更常用 JSON,而不是 Java 原生对象序列化。
1package com.example.activemq.config; 2 3import org.springframework.context.annotation.Bean; 4import org.springframework.context.annotation.Configuration; 5import org.springframework.jms.support.converter.MappingJackson2MessageConverter; 6import org.springframework.jms.support.converter.MessageConverter; 7import org.springframework.jms.support.converter.MessageType; 8 9@Configuration 10public class JmsMessageConfig { 11 12 @Bean 13 public MessageConverter jacksonJmsMessageConverter() { 14 MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter(); 15 converter.setTargetType(MessageType.TEXT); 16 converter.setTypeIdPropertyName("_type"); 17 return converter; 18 } 19} 20
Queue 实战:订单消息
常量:
1package com.example.activemq.config; 2 3public class ActiveMqDestinations { 4 5 public static final String ORDER_QUEUE = "order.queue"; 6 public static final String NOTICE_TOPIC = "notice.topic"; 7 8 private ActiveMqDestinations() { 9 } 10} 11
生产者:
1package com.example.activemq.order; 2 3import com.example.activemq.config.ActiveMqDestinations; 4import org.springframework.jms.core.JmsTemplate; 5import org.springframework.stereotype.Service; 6 7@Service 8public class OrderMessageProducer { 9 10 private final JmsTemplate jmsTemplate; 11 12 public OrderMessageProducer(JmsTemplate jmsTemplate) { 13 this.jmsTemplate = jmsTemplate; 14 } 15 16 public void sendOrderCreated(OrderCreatedMessage message) { 17 jmsTemplate.convertAndSend(ActiveMqDestinations.ORDER_QUEUE, message); 18 } 19} 20
消费者:
1package com.example.activemq.order; 2 3import com.example.activemq.config.ActiveMqDestinations; 4import org.springframework.jms.annotation.JmsListener; 5import org.springframework.stereotype.Component; 6 7@Component 8public class OrderMessageConsumer { 9 10 @JmsListener(destination = ActiveMqDestinations.ORDER_QUEUE) 11 public void handle(OrderCreatedMessage message) { 12 System.out.println("收到订单消息:" + message); 13 14 // 执行业务逻辑:扣库存、发通知、写日志等 15 } 16} 17
Controller:
1package com.example.activemq.order; 2 3import org.springframework.web.bind.annotation.PostMapping; 4import org.springframework.web.bind.annotation.RequestMapping; 5import org.springframework.web.bind.annotation.RestController; 6 7import java.math.BigDecimal; 8import java.time.LocalDateTime; 9 10@RestController 11@RequestMapping("/api/orders") 12public class OrderController { 13 14 private final OrderMessageProducer producer; 15 16 public OrderController(OrderMessageProducer producer) { 17 this.producer = producer; 18 } 19 20 @PostMapping 21 public String createOrder() { 22 OrderCreatedMessage message = new OrderCreatedMessage( 23 1001L, 24 "O202606300001", 25 2001L, 26 new BigDecimal("99.80"), 27 LocalDateTime.now() 28 ); 29 30 producer.sendOrderCreated(message); 31 return "订单消息已发送"; 32 } 33} 34
访问:
1POST http://localhost:8080/api/orders 2
同时支持 Queue 和 Topic
JmsTemplate 有一个关键属性:
1pubSubDomain 2
含义:
1false:Queue 模式 2true:Topic 模式 3
如果一个项目同时使用 Queue 和 Topic,可以分别创建两个模板。
1package com.example.activemq.config; 2 3import jakarta.jms.ConnectionFactory; 4import org.springframework.context.annotation.Bean; 5import org.springframework.context.annotation.Configuration; 6import org.springframework.jms.config.DefaultJmsListenerContainerFactory; 7import org.springframework.jms.core.JmsTemplate; 8 9@Configuration 10public class JmsTemplateConfig { 11 12 @Bean 13 public JmsTemplate queueJmsTemplate(ConnectionFactory connectionFactory) { 14 JmsTemplate template = new JmsTemplate(connectionFactory); 15 template.setPubSubDomain(false); 16 return template; 17 } 18 19 @Bean 20 public JmsTemplate topicJmsTemplate(ConnectionFactory connectionFactory) { 21 JmsTemplate template = new JmsTemplate(connectionFactory); 22 template.setPubSubDomain(true); 23 return template; 24 } 25 26 @Bean 27 public DefaultJmsListenerContainerFactory topicListenerContainerFactory( 28 ConnectionFactory connectionFactory 29 ) { 30 DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); 31 factory.setConnectionFactory(connectionFactory); 32 factory.setPubSubDomain(true); 33 return factory; 34 } 35} 36
Topic 实战:广播通知
生产者:
1package com.example.activemq.notice; 2 3import com.example.activemq.config.ActiveMqDestinations; 4import org.springframework.beans.factory.annotation.Qualifier; 5import org.springframework.jms.core.JmsTemplate; 6import org.springframework.stereotype.Service; 7 8@Service 9public class NoticeProducer { 10 11 private final JmsTemplate topicJmsTemplate; 12 13 public NoticeProducer(@Qualifier("topicJmsTemplate") JmsTemplate topicJmsTemplate) { 14 this.topicJmsTemplate = topicJmsTemplate; 15 } 16 17 public void publish(String content) { 18 topicJmsTemplate.convertAndSend(ActiveMqDestinations.NOTICE_TOPIC, content); 19 } 20} 21
短信消费者:
1package com.example.activemq.notice; 2 3import com.example.activemq.config.ActiveMqDestinations; 4import org.springframework.jms.annotation.JmsListener; 5import org.springframework.stereotype.Component; 6 7@Component 8public class SmsNoticeConsumer { 9 10 @JmsListener( 11 destination = ActiveMqDestinations.NOTICE_TOPIC, 12 containerFactory = "topicListenerContainerFactory" 13 ) 14 public void handle(String message) { 15 System.out.println("短信服务收到通知:" + message); 16 } 17} 18
邮件消费者:
1package com.example.activemq.notice; 2 3import com.example.activemq.config.ActiveMqDestinations; 4import org.springframework.jms.annotation.JmsListener; 5import org.springframework.stereotype.Component; 6 7@Component 8public class EmailNoticeConsumer { 9 10 @JmsListener( 11 destination = ActiveMqDestinations.NOTICE_TOPIC, 12 containerFactory = "topicListenerContainerFactory" 13 ) 14 public void handle(String message) { 15 System.out.println("邮件服务收到通知:" + message); 16 } 17} 18
一条 Topic 消息会被两个消费者都收到。
Topic 持久订阅 Demo
非持久订阅只接收在线期间的 Topic 消息。
持久订阅需要设置 clientId 和订阅名。
1@Bean 2public DefaultJmsListenerContainerFactory durableTopicListenerContainerFactory( 3 ConnectionFactory connectionFactory 4) { 5 DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); 6 factory.setConnectionFactory(connectionFactory); 7 factory.setPubSubDomain(true); 8 factory.setSubscriptionDurable(true); 9 factory.setClientId("notice-service-01"); 10 return factory; 11} 12
监听:
1@JmsListener( 2 destination = ActiveMqDestinations.NOTICE_TOPIC, 3 subscription = "notice-service-subscription", 4 containerFactory = "durableTopicListenerContainerFactory" 5) 6public void handleDurableTopic(String message) { 7 System.out.println("持久订阅收到通知:" + message); 8} 9
持久订阅注意点:
1同一个 broker 上 clientId 需要唯一 2同一个 clientId + subscription 代表同一个持久订阅 3
消息选择器
消息选择器可以按消息属性过滤。
发送:
1jmsTemplate.convertAndSend("order.queue", message, jmsMessage -> { 2 jmsMessage.setStringProperty("bizType", "order"); 3 jmsMessage.setStringProperty("level", "important"); 4 return jmsMessage; 5}); 6
消费:
1@JmsListener( 2 destination = "order.queue", 3 selector = "bizType = 'order' AND level = 'important'" 4) 5public void handleImportantOrder(OrderCreatedMessage message) { 6 System.out.println("收到重要订单消息:" + message); 7} 8
选择器适合按少量稳定属性过滤。
如果路由规则很复杂,通常更适合拆分队列或 Topic。
延迟消息
ActiveMQ Classic 支持调度消息。
常见属性:
| 属性 | 作用 |
|---|---|
| AMQ_SCHEDULED_DELAY | 延迟多久投递 |
| AMQ_SCHEDULED_PERIOD | 重复投递周期 |
| AMQ_SCHEDULED_REPEAT | 重复次数 |
| AMQ_SCHEDULED_CRON | Cron 表达式 |
发送延迟消息:
1import org.apache.activemq.ScheduledMessage; 2 3jmsTemplate.convertAndSend("order.timeout.queue", message, jmsMessage -> { 4 jmsMessage.setLongProperty(ScheduledMessage.AMQ_SCHEDULED_DELAY, 30 * 60 * 1000L); 5 return jmsMessage; 6}); 7
表示 30 分钟后投递。
Broker 需要开启 scheduler 支持。
死信队列
ActiveMQ Classic 默认有死信队列。
常见默认名称:
1ActiveMQ.DLQ 2
消息可能进入死信队列的常见原因:
- 消费者重复失败
- 消息过期
- 消息被回滚多次
- 达到最大重投递次数
ActiveMQ 常见处理链路:
1消费者处理失败 2 | 3 v 4Session rollback 或监听器抛异常 5 | 6 v 7Broker 重新投递 8 | 9 v 10超过最大重投递次数 11 | 12 v 13进入 DLQ 14
死信队列适合:
- 保留异常消息
- 后台补偿处理
- 人工排查
- 告警通知
Redelivery 重投递
ActiveMQ 客户端可以设置重投递策略。
1import org.apache.activemq.ActiveMQConnectionFactory; 2import org.apache.activemq.RedeliveryPolicy; 3 4ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory("tcp://localhost:61616"); 5 6RedeliveryPolicy policy = new RedeliveryPolicy(); 7policy.setMaximumRedeliveries(3); 8policy.setInitialRedeliveryDelay(1000); 9policy.setRedeliveryDelay(2000); 10policy.setUseExponentialBackOff(true); 11policy.setBackOffMultiplier(2.0); 12 13factory.setRedeliveryPolicy(policy); 14
含义:
| 配置 | 说明 |
|---|---|
| maximumRedeliveries | 最大重投递次数 |
| initialRedeliveryDelay | 首次重投递延迟 |
| redeliveryDelay | 重投递间隔 |
| useExponentialBackOff | 是否指数退避 |
| backOffMultiplier | 退避倍数 |
重投递次数耗尽后,消息通常会进入 DLQ。
Spring 事务消费
Spring JMS 可以使用事务监听容器。
1@Bean 2public DefaultJmsListenerContainerFactory txQueueListenerContainerFactory( 3 ConnectionFactory connectionFactory 4) { 5 DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); 6 factory.setConnectionFactory(connectionFactory); 7 factory.setSessionTransacted(true); 8 return factory; 9} 10
监听:
1@JmsListener( 2 destination = ActiveMqDestinations.ORDER_QUEUE, 3 containerFactory = "txQueueListenerContainerFactory" 4) 5public void handleWithTx(OrderCreatedMessage message) { 6 System.out.println("事务消费消息:" + message); 7 8 // 方法正常结束:提交 Session 9 // 方法抛出异常:回滚 Session,消息等待重投递 10} 11
事务消费适合:
- 消息处理成功才确认
- 失败后触发重投递
- 和本地数据库事务配合
幂等处理
消息中间件通常更接近“至少一次”投递。
因此,消费者需要能处理重复消息。
常见做法:
1消息带唯一 messageId 2消费前检查处理记录 3未处理则执行业务 4处理成功后保存记录 5重复消息直接跳过 6
示例:
1@JmsListener(destination = ActiveMqDestinations.ORDER_QUEUE) 2public void handle(OrderCreatedMessage message) { 3 String messageId = message.orderNo(); 4 5 if (messageLogRepository.exists(messageId)) { 6 return; 7 } 8 9 processOrder(message); 10 messageLogRepository.save(messageId); 11} 12
可以用来做幂等的方式:
- 数据库唯一索引
- 消息消费日志表
- Redis
SETNX - 业务状态机
ActiveMQ 数据存储在哪里
ActiveMQ 数据存储在 Broker 节点上。
主要包括:
- 队列和主题元数据
- 持久化消息
- 订阅关系
- 事务和调度消息信息
ActiveMQ Classic 常见持久化存储是:
1KahaDB 2
Artemis 常见持久化方式是:
1Journal 2
Docker 部署时,需要挂载数据目录。
Classic 常见目录:
1/opt/apache-activemq/data 2
不同镜像目录可能不同,需要以镜像说明为准。
示例:
1docker run -d \ 2 --name activemq-classic \ 3 -p 61616:61616 \ 4 -p 8161:8161 \ 5 -v activemq-data:/opt/apache-activemq/data \ 6 apache/activemq-classic:latest 7
如果没有挂载数据卷,容器删除后数据也会删除。
消息重启后是否保留,取决于:
1消息是否持久化 2目的地和 Broker 持久化配置 3存储目录是否保留 4
管理后台常看指标
ActiveMQ 管理后台常见关注点:
| 指标 | 说明 |
|---|---|
| Pending Messages | 等待消费的消息 |
| Consumers | 消费者数量 |
| Enqueued | 入队消息总数 |
| Dequeued | 出队消息总数 |
| Dispatch Queue | 已分派但未确认消息 |
| Expired | 过期消息 |
| DLQ | 死信消息 |
排查思路:
1Pending Messages 持续增加:消费能力不足或消费者异常 2Consumers 为 0:消费者未启动或监听配置错误 3DLQ 增加:业务处理失败、重投递耗尽或消息过期 4Enqueued 高于 Dequeued:生产速度超过消费速度 5
常见使用建议
Queue 和 Topic 分清楚
Queue 是任务分发。
1一条消息只给一个消费者 2
Topic 是广播事件。
1一条消息给多个订阅者 2
如果业务是“发短信任务”,通常选 Queue。
如果业务是“订单已支付事件”,多个系统都要收到,通常选 Topic。
JSON 优先于 Java 原生序列化
跨系统消息更适合 JSON。
原因:
- 可读性更好
- 跨语言更容易
- 版本演进更可控
- 排查消息更方便
关键消息设置持久化
发送端可以设置:
1jmsTemplate.setDeliveryPersistent(true); 2
原生 JMS 可以设置:
1producer.setDeliveryMode(DeliveryMode.PERSISTENT); 2
还需要确保 Broker 持久化存储配置正常。
消费者需要幂等
消息可能重复投递。
消费者处理逻辑需要支持重复执行。
例如:
1订单已处理 -> 直接跳过 2订单未处理 -> 执行业务并记录处理状态 3
失败消息进入 DLQ
失败消息不适合无限重试。
更常见的做法:
1有限次数重投递 2超过次数进入 DLQ 3后台补偿或人工处理 4
常用 API 汇总
| API / 注解 | 作用 |
|---|---|
| JmsTemplate.convertAndSend(...) | 发送消息 |
| @JmsListener | 监听消息 |
| ConnectionFactory | 创建连接 |
| Connection | JMS 连接 |
| Session | JMS 会话 |
| Queue | 点对点目的地 |
| Topic | 发布订阅目的地 |
| MessageProducer | 消息生产者 |
| MessageConsumer | 消息消费者 |
| TextMessage | 文本消息 |
| MappingJackson2MessageConverter | JSON 消息转换 |
| Session.CLIENT_ACKNOWLEDGE | 客户端确认 |
| session.commit() | 提交事务会话 |
| session.rollback() | 回滚事务会话 |
| RedeliveryPolicy | 重投递策略 |
总结
ActiveMQ 是 Java 生态里经典的 JMS 消息中间件。
它的核心模型很清晰:
1Queue:点对点,一条消息一个消费者 2Topic:发布订阅,一条消息多个订阅者 3
日常 Spring Boot 项目里,常用组合是:
1spring-boot-starter-activemq 2JmsTemplate 3@JmsListener 4MappingJackson2MessageConverter 5
真正落到工程实践,需要重点关注:
- Queue 和 Topic 模型选择
- JSON 消息格式
- 消息持久化
- 事务或确认机制
- 重投递策略
- DLQ 死信队列
- 消费幂等
- 管理后台监控
ActiveMQ 比 RabbitMQ 更贴近 JMS 规范。
维护传统 Java 企业系统、对接 JMS 老系统、使用 Queue / Topic 这种经典模型时,ActiveMQ 依然有明确价值。
《Java ActiveMQ 实战指南:从 JMS 队列、主题到 Spring Boot 消息处理》 是转载文章,点击查看原文。