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

Spring Integration跨Spring Boot应用JMS消息发送失败求助

Spring Integration + 嵌入式ActiveMQ:应用启动后立即终止,无消息交互问题

问题场景

尝试通过Spring Integration在两个Spring Boot应用(master/worker)之间实现JMS消息收发,使用嵌入式ActiveMQ作为代理。预期效果:

  • Worker启动后持续监听消息,接收后打印内容
  • Master启动后发送3条消息,完成后退出

实际问题:两个应用启动后立即关闭,无任何消息交互行为。

现有配置

MasterConfig

@Profile("master")
@EnableBatchProcessing
@EnableBatchIntegration
public class MasterConfig {

    @Autowired
    private JmsTemplate jmsTemplate;

    @Bean
    public MessageChannel inputChannel(){
        return MessageChannels.direct().get();
    }

    @Bean
    public IntegrationFlow jmsOutboundFlow(JmsTemplate jmsTemplate){
        return IntegrationFlows.from(inputChannel())
                .handle(Jms.outboundAdapter(jmsTemplate).destination("requestQueue"))
                .get();
    }

    @Bean
    public ApplicationRunner runner(){
        return args -> {
            for(int i = 1; i <= 3; i++)
                inputChannel().send(MessageBuilder.withPayload("hello " + i).build());
        };
    }

}

WorkerConfig

@Profile("worker")
@EnableBatchIntegration
@EnableBatchProcessing
public class WorkerConfig {

    @Autowired
    private CachingConnectionFactory cachingConnectionFactory;

    @Bean
    public IntegrationFlow jmsMessageDrivenFlow(){
        return IntegrationFlows.from(Jms.messageDrivenChannelAdapter(cachingConnectionFactory).destination("requestQueue"))
                .handle(new MessageHandler() {
                    @Override
                    public void handleMessage(Message<?> message) throws MessagingException {
                        System.out.println(message.getPayload());
                    }
                })
                .get();
    }
}

CommonConfig

@Profile("worker | master")
@Configuration
@EnableIntegration
public class CommonConfig {

    @Value("${activemq.broker-url}")
    private String brokerUrl;

    @Bean
    public ConnectionFactory connectionFactory(){
        ActiveMQConnectionFactory cf = new ActiveMQConnectionFactory();
        cf.setBrokerURL(brokerUrl);
        return cf;
    }

    @Bean
    public CachingConnectionFactory cachingConnectionFactory(){
        return new CachingConnectionFactory(connectionFactory());
    }

    @Bean
    public JmsTemplate jmsTemplate() {
        JmsTemplate jmsTemplate =
                new JmsTemplate(cachingConnectionFactory());
        jmsTemplate.setPubSubDomain(false);
        return jmsTemplate;
    }

    @Bean
    public Queue requestQueue(){
        return new ActiveMQQueue("queue.demo");
    }

    @Bean
    public Queue replyQueue(){
        return new ActiveMQQueue("queue.reply");
    }

}

pom.xml依赖

<dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-activemq</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-batch</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-integration</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.integration</groupId>
            <artifactId>spring-integration-jms</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.batch</groupId>
            <artifactId>spring-batch-integration</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.integration</groupId>
            <artifactId>spring-integration-file</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.integration</groupId>
            <artifactId>spring-integration-ftp</artifactId>
        </dependency>
        <dependency>
            <groupId>com.h2database</groupId>
            <artifactId>h2</artifactId>
            <scope>runtime</scope>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.springframework.batch</groupId>
            <artifactId>spring-batch-test</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.springframework.integration</groupId>
            <artifactId>spring-integration-test</artifactId>
            <scope>test</scope>
        </dependency>
    </dependencies>

初始application.properties

activemq.broker-url=vm://localhost

启动命令

  1. Worker端:
java -jar target/test-jms-integration-0.0.1-SNAPSHOT.jar --spring.profiles.active=worker
  1. Master端:
java -jar target/test-jms-integration-0.0.1-SNAPSHOT.jar --spring.profiles.active=master

已尝试的修改

  1. 更新application.properties:
activemq.broker-url=tcp://localhost:61616
  1. 在CommonConfig中添加BrokerService Bean:
@Bean
public BrokerService broker() throws Exception {
    BrokerService broker = new BrokerService();
    broker.addConnector("tcp://localhost:61616");
    return broker;
}
  1. pom.xml新增依赖:
<dependency>
    <groupId>org.apache.activemq</groupId>
    <artifactId>activemq-kahadb-store</artifactId>
</dependency>

解决方案

1. 核心问题定位

应用启动后立即关闭的本质是:没有非守护线程维持应用运行。同时存在队列名称不匹配的问题,导致消息无法正确路由。

2. 具体修复步骤

(1)修正Worker端运行逻辑

Worker端不需要Spring Batch相关配置,移除@EnableBatchProcessing和@EnableBatchIntegration,避免Batch模块触发应用自动退出:

@Profile("worker")
@Configuration
public class WorkerConfig {
    // 原有代码保持不变
}

或者在worker的application.properties中添加:

spring.batch.job.enabled=false

(2)统一队列名称

现有配置中,Master发送到字符串"requestQueue",但CommonConfig定义的队列Bean实际名称是queue.demo,导致消息路由错误。修改MasterConfig的出站流,注入队列Bean:

@Bean
public IntegrationFlow jmsOutboundFlow(JmsTemplate jmsTemplate, Queue requestQueue){
    return IntegrationFlows.from(inputChannel())
            .handle(Jms.outboundAdapter(jmsTemplate).destination(requestQueue))
            .get();
}

(3)正确配置嵌入式Broker

仅在Master端启动Broker,避免两个应用同时启动导致冲突:

@Profile("master")
@Bean
public BrokerService broker() throws Exception {
    BrokerService broker = new BrokerService();
    broker.addConnector("tcp://localhost:61616");
    broker.setPersistent(false); // 测试场景禁用持久化
    broker.start(); // 主动启动Broker
    return broker;
}

(4)Master端等待消息发送完成

在Master的ApplicationRunner中添加短暂等待,确保消息发送到Broker后再退出:

@Bean
public ApplicationRunner runner(MessageChannel inputChannel){
    return args -> {
        for(int i = 1; i <= 3; i++){
            inputChannel.send(MessageBuilder.withPayload("hello " + i).build());
            System.out.println("已发送消息:hello " + i);
        }
        Thread.sleep(2000); // 等待消息投递完成
    };
}

(5)清理不必要依赖

移除未使用的Spring Batch、文件/FTP集成依赖,减少环境干扰:

<!-- 移除以下依赖 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-batch</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.batch</groupId>
    <artifactId>spring-batch-integration</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-file</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-ftp</artifactId>
</dependency>

3. 验证流程

  1. 先启动Worker应用,确认应用保持运行状态
  2. 启动Master应用,观察控制台输出“已发送消息”日志
  3. 查看Worker控制台,确认打印出hello 1、hello 2、hello 3

内容的提问来源于stack exchange,提问作者of32 inc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 00:02:01