求基于Java Config的Spring Integration与JMS业务场景配置示例
嘿,作为Spring Integration新手,你的这个业务场景刚好能用上它的核心组件来解耦逻辑,替代原来硬编码在@JmsListener里的实现。我给你整理了一套完整的Java配置方案,完全覆盖你提到的所有需求:从WebSphere MQ取消息、按消息头路由、存MongoDB、错误信息单独存储。
整体架构思路
我们会用到这些Spring Integration核心组件:
- JMS入站通道适配器:替代
@JmsListener,从WebSphere MQ队列接收消息并发送到Integration通道 - 消息头路由器:根据指定的消息头值,把消息路由到对应业务服务的通道
- 服务激活器:绑定业务服务方法,处理消息并存储到MongoDB主集合
- 错误处理通道+处理器:捕获整个流程中的异常,把错误信息存储到MongoDB的错误集合
1. 依赖准备
先确保你的pom.xml里包含这些必要依赖(如果是Gradle的话对应转换即可):
<dependencies> <!-- Spring Boot核心依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-integration</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-jms</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-mongodb</artifactId> </dependency> <!-- WebSphere MQ JMS驱动 --> <dependency> <groupId>com.ibm.mq</groupId> <artifactId>mq-jms-spring-boot-starter</artifactId> <version>2.7.4</version> <!-- 用最新稳定版即可 --> </dependency> </dependencies>
2. Java配置实现
我们把配置拆成几个部分,方便理解和维护:
2.1 JMS入站适配器配置
这个组件负责从WebSphere MQ接收消息,替换原来的@JmsListener:
@Configuration @EnableIntegration public class IntegrationConfig { // WebSphere MQ连接工厂(Spring Boot会自动绑定application.properties里的配置) @Autowired private ConnectionFactory mqConnectionFactory; // 定义输入通道:MQ消息会先进入这个通道 @Bean public MessageChannel mqInputChannel() { return new DirectChannel(); } // JMS消息驱动通道适配器:绑定MQ队列和输入通道 @Bean public JmsMessageDrivenChannelAdapter jmsMessageDrivenChannelAdapter() { JmsMessageDrivenChannelAdapter adapter = new JmsMessageDrivenChannelAdapter(mqConnectionFactory); adapter.setDestinationName("YOUR_MQ_QUEUE_NAME"); // 替换成你的MQ队列名 adapter.setOutputChannel(mqInputChannel()); // 可选:配置并发消费者数量 adapter.setConcurrentConsumers(3); return adapter; } }
2.2 消息头路由配置
这里我们根据消息头(比如叫serviceRoute)的值,把消息路由到不同的业务服务通道:
@Configuration public class IntegrationConfig { // ... 前面的JMS配置 ... // 定义各个业务服务的通道 @Bean public MessageChannel service1Channel() { return new DirectChannel(); } @Bean public MessageChannel service2Channel() { return new DirectChannel(); } @Bean public MessageChannel defaultServiceChannel() { return new DirectChannel(); } // 消息头路由器:根据serviceRoute头的值路由 @Bean public HeaderValueRouter headerValueRouter() { HeaderValueRouter router = new HeaderValueRouter("serviceRoute"); router.setChannelMapping("SERVICE_1", "service1Channel"); router.setChannelMapping("SERVICE_2", "service2Channel"); router.setDefaultOutputChannel("defaultServiceChannel"); // 处理未知头值的情况 router.setInputChannel(mqInputChannel()); // 绑定到MQ输入通道 return router; } }
2.3 业务服务与MongoDB存储
编写业务服务类,处理消息并存储到MongoDB的business_messages集合:
@Service public class BusinessMessageService { private static final Logger log = LoggerFactory.getLogger(BusinessMessageService.class); @Autowired private MongoTemplate mongoTemplate; // 处理SERVICE_1的消息 @ServiceActivator(inputChannel = "service1Channel") public void handleService1Message(Message<String> message) { // 这里写你的SERVICE_1业务逻辑 String payload = message.getPayload(); // 构造存储对象 BusinessMessage businessMessage = new BusinessMessage(); businessMessage.setPayload(payload); businessMessage.setRouteKey("SERVICE_1"); businessMessage.setReceivedTime(LocalDateTime.now()); // 存储到MongoDB mongoTemplate.save(businessMessage, "business_messages"); } // 处理SERVICE_2的消息 @ServiceActivator(inputChannel = "service2Channel") public void handleService2Message(Message<String> message) { // 这里写你的SERVICE_2业务逻辑 String payload = message.getPayload(); BusinessMessage businessMessage = new BusinessMessage(); businessMessage.setPayload(payload); businessMessage.setRouteKey("SERVICE_2"); businessMessage.setReceivedTime(LocalDateTime.now()); mongoTemplate.save(businessMessage, "business_messages"); } // 处理未知路由的消息 @ServiceActivator(inputChannel = "defaultServiceChannel") public void handleDefaultMessage(Message<String> message) { // 可以记录日志或者做降级处理 log.warn("Received message with unknown route key: {}", message.getHeaders().get("serviceRoute")); // 也可以存储到MongoDB BusinessMessage businessMessage = new BusinessMessage(); businessMessage.setPayload(message.getPayload()); businessMessage.setRouteKey("DEFAULT"); businessMessage.setReceivedTime(LocalDateTime.now()); mongoTemplate.save(businessMessage, "business_messages"); } // 对应的实体类 public static class BusinessMessage { private String id; private String payload; private String routeKey; private LocalDateTime receivedTime; // getter和setter省略 } }
2.4 错误处理配置
捕获整个流程中的异常,把错误信息存储到MongoDB的error_messages集合:
@Service public class ErrorHandlingService { @Autowired private MongoTemplate mongoTemplate; // 绑定到全局错误通道 @ServiceActivator(inputChannel = "errorChannel") public void handleError(ErrorMessage errorMessage) { // 获取原消息和异常信息 Message<?> originalMessage = errorMessage.getOriginalMessage(); Throwable exception = errorMessage.getPayload(); // 构造错误存储对象 ErrorLog errorLog = new ErrorLog(); errorLog.setOriginalPayload(originalMessage != null ? originalMessage.getPayload().toString() : null); errorLog.setErrorReason(exception.getMessage()); errorLog.setStackTrace(ExceptionUtils.getStackTrace(exception)); // 需要导入org.apache.commons.lang3.ExceptionUtils errorLog.setErrorTime(LocalDateTime.now()); // 存储到错误集合 mongoTemplate.save(errorLog, "error_messages"); } // 对应的错误实体类 public static class ErrorLog { private String id; private String originalPayload; private String errorReason; private String stackTrace; private LocalDateTime errorTime; // getter和setter省略 } }
3. 配置文件(application.properties)
添加WebSphere MQ和MongoDB的连接配置:
# WebSphere MQ配置 ibm.mq.queue-manager=YOUR_QUEUE_MANAGER ibm.mq.channel=YOUR_CHANNEL ibm.mq.conn-name=YOUR_MQ_HOST(YOUR_MQ_PORT) ibm.mq.user=YOUR_MQ_USER ibm.mq.password=YOUR_MQ_PASSWORD # MongoDB配置 spring.data.mongodb.uri=mongodb://localhost:27017/your_database_name
可选:XML配置参考
如果你需要XML配置,核心结构大概是这样:
<int:channel id="mqInputChannel"/> <int:channel id="service1Channel"/> <int:channel id="service2Channel"/> <int:channel id="defaultServiceChannel"/> <int-jms:message-driven-channel-adapter id="mqAdapter" connection-factory="mqConnectionFactory" destination-name="YOUR_MQ_QUEUE_NAME" output-channel="mqInputChannel"/> <int:header-value-router input-channel="mqInputChannel" header-name="serviceRoute"> <int:mapping value="SERVICE_1" channel="service1Channel"/> <int:mapping value="SERVICE_2" channel="service2Channel"/> <int:default-output-channel ref="defaultServiceChannel"/> </int:header-value-router> <int:service-activator input-channel="service1Channel" ref="businessMessageService" method="handleService1Message"/> <int:service-activator input-channel="service2Channel" ref="businessMessageService" method="handleService2Message"/> <int:service-activator input-channel="defaultServiceChannel" ref="businessMessageService" method="handleDefaultMessage"/> <int:service-activator input-channel="errorChannel" ref="errorHandlingService" method="handleError"/>
这样整个流程就完全解耦了:MQ接收、路由、业务处理、存储、错误处理都拆成了独立组件,后续修改逻辑只需要调整对应的部分就行。要是你有细节问题,比如MQ的特殊配置、Mongo的自定义序列化,随时交流~
内容的提问来源于stack exchange,提问作者toppless
相关产品推荐
相关产品推荐

