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

Spring Integration中JDBC入站出站适配器连接及列表消息处理问题

问题描述

我想要编写一个实现表复制功能的示例应用,相关配置及问题如下:

  • 步骤1:测试jdbc:inbound-channel-adapter与jdbcMessageHandler的连接,此步骤运行正常。
  • 步骤2:尝试连接jdbc:outbound-channel-adapter与入站适配器,当前配置代码如下:
<?xml version="1.0" encoding="UTF-8"?><beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:int="http://www.springframework.org/schema/integration"
xmlns:int-jdbc="http://www.springframework.org/schema/integration/jdbc"
xsi:schemaLocation="http://www.springframework.org/schema/beans https://www.springframework.org/schema/beans/spring-beans.xsd
    http://www.springframework.org/schema/integration https://www.springframework.org/schema/integration/spring-integration.xsd
    http://www.springframework.org/schema/integration/jdbc http://www.springframework.org/schema/integration/jdbc/spring-integration-jdbc.xsd">

<int-jdbc:inbound-channel-adapter id="dataChannel"
    query="select * from articles where sent = 0" 
    update="update articles set sent = 1 where id in (:id)"
    data-source="dataSource" row-mapper="articleRowMapper">
    <int:poller fixed-rate="10000">
        <int:transactional />
    </int:poller>
</int-jdbc:inbound-channel-adapter>


<!-- <int:service-activator input-channel="dataChannel" ref="jdbcMessageHandler" /> 
<bean id="jdbcMessageHandler" class="com.demo.service.JdbcMessageHandler" /> -->


<bean id="transactionManager"
    class="org.springframework.jdbc.datasource.DataSourceTransactionManager">
    <property name="dataSource" ref="dataSource" />
</bean>

<int:poller default="true" fixed-rate="10000" />

<int:channel id="dataChannel">
    <int:queue />
</int:channel>


<bean id="dataSource"
    class="org.springframework.jdbc.datasource.DriverManagerDataSource">
    <property name="driverClassName" value="com.mysql.jdbc.Driver" />
    <property name="url" value="jdbc:mysql://localhost:3306/test" />
    <property name="username" value="demo" />
    <property name="password" value="password" />
</bean>

<bean id="articleRowMapper" class="com.demo.domain.ArticleRowMapper" />

<int-jdbc:outbound-channel-adapter id="jdbcOutbound"
                                       channel="dataChannel"
                                       data-source="dataSource"
                                       sql-parameter-source-factory="sqlParameterSource"
                                       query="INSERT INTO ARTICLES(ID, NAME, CATEGORY , TAGS , AUTHOR , SENT) VALUES(:id, :name, :category,:tags,:author,:sent)"/>
    <bean id="sqlParameterSource"   class="org.springframework.integration.jdbc.ExpressionEvaluatingSqlParameterSourceFactory">
        <property name="parameterExpressions">
            <map>
                <entry key="id" value="payload.id"/>
                <entry key="name" value="payload.name"/>
                <entry key="category" value="payload.category"/>
                <entry key="tags" value="payload.tags"/>
                <entry key="author" value="payload.author"/>
                <entry key="sent" value="payload.sent"/>                    
            </map>
        </property>
    </bean>
</beans>

运行时出现如下错误:

SEVERE: org.springframework.messaging.MessageHandlingException: error
occurred in message handler
[org.springframework.integration.jdbc.JdbcMessageHandler#0]; nested
exception is
org.springframework.dao.TransientDataAccessResourceException:
PreparedStatementCallback; SQL [INSERT INTO ARTICLES(ID, NAME,
CATEGORY , TAGS , AUTHOR , SENT) VALUES(?, ?, ?,?,?,?)Invalid argument
value: java.io.NotSerializableException; nested exception is
java.sql.SQLException: Invalid argument value:
java.io.NotSerializableException, failedMessage=GenericMessage
[payload=[com.demo.domain.Article@49023a31,
com.demo.domain.Article@2d25f418, com.demo.domain.Article@31d44fe8],
headers={id=1c16e823-b70e-faf1-2ea5-5475e3f8dbde,
timestamp=1661550341366}] at
org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:153)

我推测入站适配器返回的是List类型消息,而出站适配器仅支持处理单条消息,请问该如何处理这类List消息?此外,我希望避免使用同一个dataChannel以实现松耦合,该如何操作?

解决方案

1. 处理List类型消息

你的推测正确,jdbc:inbound-channel-adapter默认会把查询结果封装成List作为消息payload,而jdbc:outbound-channel-adapter仅能处理单条数据。解决这个问题需要使用splitter组件拆分集合,将List中的每个元素转为独立消息:

在配置中添加splitter并调整通道流向:

<!-- 入站适配器输出到独立通道 -->
<int-jdbc:inbound-channel-adapter id="jdbcInbound"
    channel="inboundChannel"
    query="select * from articles where sent = 0" 
    update="update articles set sent = 1 where id in (:id)"
    data-source="dataSource" row-mapper="articleRowMapper">
    <int:poller fixed-rate="10000">
        <int:transactional />
    </int:poller>
</int-jdbc:inbound-channel-adapter>

<!-- splitter自动拆分List,将每个Article转为单独消息,输出到splitChannel -->
<int:splitter input-channel="inboundChannel" output-channel="splitChannel"/>

<!-- 拆分后的单条消息通道 -->
<int:channel id="splitChannel"/>

<!-- 出站适配器监听splitChannel,处理单条消息 -->
<int-jdbc:outbound-channel-adapter id="jdbcOutbound"
                                   channel="splitChannel"
                                   data-source="dataSource"
                                   sql-parameter-source-factory="sqlParameterSource"
                                   query="INSERT INTO ARTICLES(ID, NAME, CATEGORY , TAGS , AUTHOR , SENT) VALUES(:id, :name, :category,:tags,:author,:sent)"/>

splitter会自动识别payload是否为集合/数组,将其拆分为多个独立消息,这样出站适配器就能逐个处理每个Article对象。

另外,错误中的NotSerializableException是因为队列通道要求消息payload可序列化:如果不需要队列通道,可直接去掉<int:queue />配置;如果必须保留队列,需要让Article类实现Serializable接口。

2. 实现松耦合(避免共用通道)

要实现松耦合,需为不同组件分配独立通道,通过中间组件(如splitter)间接连接,而非让入站和出站直接共用同一通道:

  • 为入站适配器指定独立的输出通道(inboundChannel)
  • 用splitter连接入站通道与拆分后的通道(splitChannel)
  • 出站适配器监听拆分后的通道(splitChannel)

这样每个组件仅依赖自身对应的通道,组件间通过通道间接交互,耦合度更低。完整调整后的核心配置如下:

<!-- 入站适配器专用通道 -->
<int:channel id="inboundChannel">
    <!-- 不需要队列可删除此配置,规避序列化要求 -->
    <!-- <int:queue /> -->
</int:channel>

<int-jdbc:inbound-channel-adapter id="jdbcInbound"
    channel="inboundChannel"
    query="select * from articles where sent = 0" 
    update="update articles set sent = 1 where id in (:id)"
    data-source="dataSource" row-mapper="articleRowMapper">
    <int:poller fixed-rate="10000">
        <int:transactional />
    </int:poller>
</int-jdbc:inbound-channel-adapter>

<!-- splitter拆分集合 -->
<int:splitter input-channel="inboundChannel" output-channel="splitChannel"/>

<!-- 拆分后单条消息通道 -->
<int:channel id="splitChannel"/>

<!-- 出站适配器监听拆分后的通道 -->
<int-jdbc:outbound-channel-adapter id="jdbcOutbound"
                                   channel="splitChannel"
                                   data-source="dataSource"
                                   sql-parameter-source-factory="sqlParameterSource"
                                   query="INSERT INTO ARTICLES(ID, NAME, CATEGORY , TAGS , AUTHOR , SENT) VALUES(:id, :name, :category,:tags,:author,:sent)"/>

内容的提问来源于stack exchange,提问作者Kris Swat

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 05:06:28