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

