Spring Integration Channel首次正常后续跨流调用失效问题
问题现象
在使用Spring Integration通过Channel对接其他集成流时,仅首次调用可正常执行,后续调用会直接跳过Channel返回。
实际运行表现为隔条消息才能成功透传:每次成功调用后下一次请求会跳过Channel直接返回,再下一次调用又恢复正常;若重复定义jamsSubmitJob Bean,则可连续成功2次后再次失效,循环往复。
涉及的上游集成流代码片段:
.handle((p, h) -> { System.out.println("Payload Before Channel" + p.toString()); return p; }) .channel(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.getChannelName()) .handle((p, h) -> { System.out.println("Payload After Channel" + p.toString()); return p; })
涉及的下游对接集成流代码:
@Bean public IntegrationFlow jamsSubmitJob() { return IntegrationFlows.from(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.getChannelName()) .handle((p, h) -> { try { jamsToken = authMang.getJamsAuth().getTokenWithTokenType(); } catch (Exception e) { e.printStackTrace(); } JAMS_SUBMIT_JOB_INTGRTN .info(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.toString() + " Integration called."); JAMS_SUBMIT_JOB_INTGRTN.info(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.toString() + " Headers:= " + h); JAMS_SUBMIT_JOB_INTGRTN.debug(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.toString() + " Payload:= " + p); return p; }) .handle((p, h) -> { // hail mary to get new token return MessageBuilder .withPayload(p) .removeHeaders("*") .setHeader(HttpHeaders.AUTHORIZATION.toLowerCase(), jamsToken) .setHeader(HttpHeaders.CONTENT_TYPE.toLowerCase(), "application/json") .build(); }) .handle((p, h) -> { JAMS_SUBMIT_JOB_INTGRTN .info(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.toString() + " Submitting payload to JAMS:"); JAMS_SUBMIT_JOB_INTGRTN.info(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.toString() + " Headers:= " + h); JAMS_SUBMIT_JOB_INTGRTN.info(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.toString() + " Payload:= " + p); return p; }) .handle(Http.outboundGateway(JAMS_SUBMIT_ENDPOINT) .requestFactory(alliantPooledHttpConnection.get_httpComponentsClientHttpRequestFactory()) .httpMethod(HttpMethod.POST) .expectedResponseType(String.class) .extractPayload(true)) .logAndReply(); }
问题根因
核心原因是Spring Integration Java DSL的通道订阅规则和默认通道分发策略不匹配:
- 上游流在
.channel(channelName)节点后直接追加.handle()处理逻辑时,框架会自动将这个后续handle注册为该命名Channel的订阅消费者,和通过IntegrationFlows.from(channelName)定义的下游jamsSubmitJob流是完全对等的消费者关系。 - 直接通过字符串名称引用、未显式指定类型的Channel,框架默认实例化为
DirectChannel类型,该类型采用轮询(Round-Robin)点对点分发策略,会把投递到通道的消息按顺序轮流发给所有订阅的消费者,不会广播。
当前场景下该Channel的消费者列表为:
- 上游流中channel节点之后的handle逻辑
- 自定义的jamsSubmitJob下游集成流
第一次消息被轮询分发给jamsSubmitJob流,执行完正常返回;第二次消息轮询分发给上游流自身的后置handle,根本没有进入下游逻辑,表现为跳过Channel直接返回;第三次消息再次轮询到jamsSubmitJob流,又恢复正常。如果重复定义2个jamsSubmitJob Bean,该Channel的消费者总数变为3个,就会连续2次将消息分发给下游jamsSubmitJob流,第三次分发给上游后置handle,和观测到的现象完全吻合。
修复方案
- 优先方案:如果预期是消息投递到目标节点后,等待下游流执行完成再回到上游流继续执行后续逻辑,将上游流中的
.channel()调用替换为.gateway()即可。.gateway()不会给目标Channel注册额外消费者,而是以请求/响应模式向通道发送消息,等待下游返回结果后再继续当前流的后续处理,完全符合透传预期。修改后的上游流代码:.handle((p, h) -> { System.out.println("Payload Before Channel" + p.toString()); return p; }) .gateway(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.getChannelName()) .handle((p, h) -> { System.out.println("Payload After Channel" + p.toString()); return p; }) - 可选方案:如果不需要下游处理结果回到上游流,直接删除上游流中
.channel()节点之后的所有处理逻辑,避免上游流成为目标Channel的消费者,防止轮询策略将消息分发给无关的处理节点。 - 注意事项:不要重复定义同名的IntegrationFlow Bean,否则会导致同一通道被多个重复流订阅,引发轮询异常、重复消费、请求响应匹配错乱等问题。
内容的提问来源于stack exchange,提问作者Jim Dannucci
相关产品推荐
相关产品推荐

