Spring Cloud Dataflow中构建数据库转JSON的IntegrationFlow报错求助
我在Spring Cloud Dataflow环境中实现从数据库读取记录并转换为JSON格式的功能,但构建IntegrationFlow时遇到了如下错误:
Caused by: org.springframework.messaging.core.DestinationResolutionException: no output-channel or replyChannel header available
at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutput(AbstractMessageProducingHandler.java:440)
at org.springframework.integration.handler.AbstractMessageProducingHandler.doProduceOutput(AbstractMessageProducingHandler.java:319)
at org.springframework.integration.handler.AbstractMessageProducingHandler.produceOutput(AbstractMessageProducingHandler.java:267)
at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutputs(AbstractMessageProducingHandler.java:231)
at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:140)
at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:62)
我的相关配置代码如下:
@Bean public MessageSource<Object> jdbcMessageSource() { String query = "select cd_int_controle, de_tabela from int_controle rowlock readpast " + "where id_status = 0 order by cd_int_controle"; JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(dataSource, query); adapter.setMaxRows(properties.getPollSize()); adapter.setUpdatePerRow(true); adapter.setRowMapper((RowMapper<IntControle>) (rs, i) -> new IntControle(rs.getLong(1), rs.getString(2))); adapter.setUpdateSql("update int_controle set id_status = 1 where cd_int_controle = :cdIntControle"); return adapter; } @Bean public IntegrationFlow jsonSupplier() { return IntegrationFlows.from(jdbcMessageSource(), c -> c.poller(Pollers.fixedRate(properties.getPollRateMs(), TimeUnit.MILLISECONDS).transactional())) .transform((GenericTransformer<List<IntControle>, String>) ints -> { //transform to Json }) .get(); }
我怀疑是IntegrationFlow中缺少必要的调用,希望能得到帮助排查问题。
问题原因与修复方案
你遇到的这个错误本质上是Spring Integration在“抱怨”:它完成了数据库查询和JSON转换的操作,但找不到后续要发送消息的目标通道,所以抛出了DestinationResolutionException。
具体怎么修?
你的IntegrationFlow在transform步骤之后缺少了消息的“归宿”,需要添加一个终端节点来指定消息的去向,结合Spring Cloud Dataflow的使用场景,给你两种常见的修复方式:
方式1:发送到Dataflow绑定的输出通道(推荐用于生产环境)
在Spring Cloud Dataflow中,我们通常会把处理后的消息发送到预定义的output通道,这个通道会自动绑定到你部署应用时指定的消息中间件目标(比如Kafka Topic、RabbitMQ队列):
@Bean public IntegrationFlow jsonSupplier() { return IntegrationFlows.from(jdbcMessageSource(), c -> c.poller(Pollers.fixedRate(properties.getPollRateMs(), TimeUnit.MILLISECONDS).transactional())) .transform((GenericTransformer<List<IntControle>, String>) ints -> { // 补全你的JSON转换逻辑,比如用Jackson的ObjectMapper ObjectMapper objectMapper = new ObjectMapper(); try { return objectMapper.writeValueAsString(ints); } catch (JsonProcessingException e) { throw new RuntimeException("转换为JSON失败", e); } }) .channel("output") // 这里指定输出通道,Dataflow会自动处理绑定 .get(); }
方式2:添加本地处理器(适合测试或调试)
如果是本地开发测试,你可以直接添加一个handle处理器来消费消息,比如打印到控制台:
@Bean public IntegrationFlow jsonSupplier() { return IntegrationFlows.from(jdbcMessageSource(), c -> c.poller(Pollers.fixedRate(properties.getPollRateMs(), TimeUnit.MILLISECONDS).transactional())) .transform((GenericTransformer<List<IntControle>, String>) ints -> { // 补全JSON转换逻辑 ObjectMapper objectMapper = new ObjectMapper(); try { return objectMapper.writeValueAsString(ints); } catch (JsonProcessingException e) { throw new RuntimeException("转换为JSON失败", e); } }) .handle(message -> { // 这里可以自定义消息处理逻辑,比如打印 System.out.println("处理后的JSON消息: " + message.getPayload()); }) .get(); }
额外提醒
- 确保你的
IntControle实体类是可序列化的(如果要发送到消息中间件),可以实现Serializable接口。 - 你的
JdbcPollingChannelAdapter配置看起来没问题,但要注意updateSql中的:cdIntControle参数要和IntControle类的cdIntControle属性名称对应,否则更新操作可能会失败。
内容的提问来源于stack exchange,提问作者Thiago Sayão

