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

求基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:02:19