You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot JMS结合ActiveMQ Classic定时发送疑似存在BUG

ActiveMQ 定时重复消息首次投递延迟问题

问题背景

根据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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 13:05:57