Spring Boot JMS结合ActiveMQ Classic定时发送疑似存在BUG
问题背景
根据Apache ActiveMQ文档,使用Spring JMS发送消息时,定时相关属性的预期行为如下:
AMQ_SCHEDULED_DELAY:消息等待代理调度投递的毫秒数AMQ_SCHEDULED_PERIOD:首次投递后再次调度的间隔毫秒数AMQ_SCHEDULED_REPEAT:重复调度投递的次数
但实际测试发现,当同时使用AMQ_SCHEDULED_PERIOD和AMQ_SCHEDULED_REPEAT时,即使显式将AMQ_SCHEDULED_DELAY设为0,首次消息投递仍会被延迟。例如设置AMQ_SCHEDULED_PERIOD=1000、AMQ_SCHEDULED_REPEAT=2,消息在时间0发送,却在1000、2000、3000毫秒时才被投递。
为验证问题,搭建了Spring Boot测试应用,结果显示标记为“delayed repeated”的消息首次接收延迟约1444毫秒,与“regular”消息的即时接收差异明显。但通过ActiveMQ Broker UI发送时,首次消息可立即投递。请问这是BUG还是对文档的理解有误?
测试代码
ExampleApplication.java
package curious.jms.delay.bug.example; import jakarta.jms.JMSException; import jakarta.jms.TextMessage; import java.time.Duration; import java.time.Instant; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; import org.apache.activemq.ScheduledMessage; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.ApplicationContext; import org.springframework.context.event.EventListener; import org.springframework.jms.annotation.EnableJms; import org.springframework.jms.annotation.JmsListener; import org.springframework.jms.core.JmsTemplate; @SpringBootApplication @EnableJms public class ExampleApplication { @Autowired JmsTemplate jmsTemplate; @Autowired private ApplicationContext context; private Map<String, Instant> messageSendTimes = new HashMap<>(); private Map<String, List<Instant>> messageReceiveTimes = new HashMap<>(); private final String QUEUE_NAME = "example.queue"; public static void main(String[] args) { SpringApplication.run(ExampleApplication.class, args); } @EventListener(ApplicationReadyEvent.class) public void run() throws InterruptedException { sendRegularMessage(); sendDelayedMessage(); sendDelayedRepeatedMessage(); Thread.sleep(10000); printResults(messageSendTimes, messageReceiveTimes); SpringApplication.exit(context, () -> 0); } private void sendRegularMessage() { String text = "regular"; System.out.println("queueing message " + text); messageSendTimes.put(text, Instant.now()); jmsTemplate.send(QUEUE_NAME, s -> s.createTextMessage(text)); } private void sendDelayedMessage() { String text = "delayed"; System.out.println("queueing message " + text); messageSendTimes.put(text, Instant.now()); jmsTemplate.send(QUEUE_NAME, s -> { TextMessage m = s.createTextMessage(text); m.setLongProperty(ScheduledMessage.AMQ_SCHEDULED_DELAY, 1000); return m; }); } private void sendDelayedRepeatedMessage() { String text = "delayed repeated"; System.out.println("queueing message " + text); messageSendTimes.put(text, Instant.now()); jmsTemplate.send(QUEUE_NAME, s -> { TextMessage m = s.createTextMessage(text); m.setLongProperty(ScheduledMessage.AMQ_SCHEDULED_DELAY, 0); m.setLongProperty(ScheduledMessage.AMQ_SCHEDULED_PERIOD, 1000); m.setIntProperty(ScheduledMessage.AMQ_SCHEDULED_REPEAT, 5); return m; }); } @JmsListener(destination = QUEUE_NAME) public void listener(TextMessage m) throws JMSException { if (messageReceiveTimes.containsKey(m.getText())) { messageReceiveTimes.get(m.getText()).add(Instant.now()); } else { List<Instant> l = new ArrayList<>(); l.add(Instant.now()); messageReceiveTimes.put(m.getText(), l); } System.out.println("received message " + m.getText() + " at " + Instant.now()); } private void printResults(Map<String, Instant> sendTimes, Map<String, List<Instant>> receiveTimes) { System.out.printf("%16s | %20s | %20s | %10s %n", "message", "sendTime", "receiveTime", "diff"); for (String k : sendTimes.keySet()) { Instant sendTime = sendTimes.get(k); for (Instant instant : receiveTimes.get(k)) { System.out.printf("%16s | %20s | %20s | %8s %n", k, formatInstant(sendTime), formatInstant(instant), Duration.between(sendTime, instant).toMillis()); } } } private String formatInstant(Instant i) { return DateTimeFormatter.ISO_TIME.format(i.atOffset(ZoneOffset.UTC)); } }
application.properties
spring.application.name=example spring.activemq.broker-url=vm://localhost?broker.persistent=true&broker.schedulerSupport=true
build.gradle
plugins { id 'java' id 'org.springframework.boot' version '3.4.4' id 'io.spring.dependency-management' version '1.1.7' } group = 'curious.jms.delay.bug' version = '0.0.1-SNAPSHOT' java { toolchain { languageVersion = JavaLanguageVersion.of(17) } } repositories { mavenCentral() } dependencies { implementation 'org.springframework.boot:spring-boot-starter-activemq' implementation 'org.apache.activemq:activemq-kahadb-store:6.1.6' testImplementation 'org.springframework.boot:spring-boot-starter-test' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' } tasks.named('test') { useJUnitPlatform() }
运行输出
message | sendTime | receiveTime | diff delayed repeated | 19:22:12.0743304Z | 19:22:13.5186246Z | 1444 delayed repeated | 19:22:12.0743304Z | 19:22:14.5030668Z | 2428 delayed repeated | 19:22:12.0743304Z | 19:22:15.5038109Z | 3429 delayed repeated | 19:22:12.0743304Z | 19:22:16.5035174Z | 4429 delayed repeated | 19:22:12.0743304Z | 19:22:17.6968498Z | 5622 delayed repeated | 19:22:12.0743304Z | 19:22:18.5176954Z | 6443 delayed | 19:22:12.0642786Z | 19:22:13.5080843Z | 1443 regular | 19:22:12.0005657Z | 19:22:12.0672787Z | 66
分析与解决方案
这并非对文档的理解错误,而是ActiveMQ客户端与Broker交互时的已知行为:
当同时设置AMQ_SCHEDULED_PERIOD和AMQ_SCHEDULED_REPEAT时,即使显式设置AMQ_SCHEDULED_DELAY=0,客户端会默认将首次投递延迟设为与AMQ_SCHEDULED_PERIOD相同的值。而通过Broker UI发送时,UI直接将调度参数传递给Broker,绕过了客户端的默认逻辑,因此首次消息可立即投递。
解决方法很简单:完全不设置AMQ_SCHEDULED_DELAY属性。当未指定该属性时,Broker会立即投递首次消息,之后按照AMQ_SCHEDULED_PERIOD的间隔重复投递AMQ_SCHEDULED_REPEAT次,完全符合文档描述的预期行为。
修改后的sendDelayedRepeatedMessage方法示例:
private void sendDelayedRepeatedMessage() { String text = "delayed repeated"; System.out.println("queueing message " + text); messageSendTimes.put(text, Instant.now()); jmsTemplate.send(QUEUE_NAME, s -> { TextMessage m = s.createTextMessage(text); // 移除AMQ_SCHEDULED_DELAY的设置 m.setLongProperty(ScheduledMessage.AMQ_SCHEDULED_PERIOD, 1000); m.setIntProperty(ScheduledMessage.AMQ_SCHEDULED_REPEAT, 5); return m; }); }
内容的提问来源于stack exchange,提问作者fantasyTeapot

