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
启动命令
- Worker端:
java -jar target/test-jms-integration-0.0.1-SNAPSHOT.jar --spring.profiles.active=worker
- Master端:
java -jar target/test-jms-integration-0.0.1-SNAPSHOT.jar --spring.profiles.active=master
已尝试的修改
- 更新application.properties:
activemq.broker-url=tcp://localhost:61616
- 在CommonConfig中添加BrokerService Bean:
@Bean public BrokerService broker() throws Exception { BrokerService broker = new BrokerService(); broker.addConnector("tcp://localhost:61616"); return broker; }
- 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. 验证流程
- 先启动Worker应用,确认应用保持运行状态
- 启动Master应用,观察控制台输出“已发送消息”日志
- 查看Worker控制台,确认打印出
hello 1、hello 2、hello 3
内容的提问来源于stack exchange,提问作者of32 inc
相关产品推荐
相关产品推荐

