Spring Integration WebFlux流程中发送Slack通知并维持主流程执行问题
1. 主流程中断原因
你在主流程中直接调用.channel(SLACK_CHANNEL)将消息送入Slack通道,而你定义的SLACK_CHANNEL是点对点的队列通道,消息一旦被送入该通道就会脱离主流程,自然不会执行后续的split、aggregate等步骤。
2. Slack发送流程未触发原因
队列通道的消费者需要配置轮询器才能主动拉取通道内的消息,你的sendFileProcessingInfo流程没有配置轮询规则,因此不会消费SLACK_CHANNEL内的消息。
3. expectedResponseType相关说明
如果你仅需要发送Slack通知无需处理返回结果,可直接使用WebFlux.outboundChannelAdapter替代outboundGateway,无需配置expectedResponseType。如果使用outboundGateway,Slack开放接口返回的是标准JSON字符串,设置为String是合理的,该配置不是强制要求,但如果设置的类型和返回值结构不匹配会触发反序列化异常。
4. prepareSlackMessage是否影响主流程消息
Spring Integration中的Message对象本身是不可变的,transform操作会返回新的Message实例,只要你没有修改原消息payload内部的可变属性(比如自定义对象的字段值),该转换操作完全不会影响主流程的消息。
你需要的是不影响主流程的旁路通知能力,Spring Integration 提供的wireTap(线路窃听)组件完全匹配该场景,它会将主流程的消息复制一份送入指定旁路通道,主流程不受任何影响继续执行。
修正后的代码实现
1. 通道定义(推荐使用异步执行通道,无需轮询)
@Bean MessageChannel slackChannel() { // 用线程池异步执行Slack发送逻辑,不阻塞主流程 return MessageChannels.executor(SLACK_CHANNEL, Executors.newCachedThreadPool()).get(); }
2. 主流程调整(移除原来的.channel(SLACK_CHANNEL),替换为wireTap)
@Bean IntegrationFlow startFlow() { return IntegrationFlows .from(FILE_CHANNEL) .filter(myFilter) .handle(myService, "doSomething") .transform(doSomeTransformation) // 新增wireTap,复制消息送入Slack通道,主流程继续执行后续步骤 .wireTap(SLACK_CHANNEL) .split() .aggregate(myAggregator) .transform(anotherTransformer) .channel(ANOTHER_CHANNEL) .get(); }
3. Slack发送流程调整(无需配置轮询,不需要返回结果时用outboundChannelAdapter)
@Bean IntegrationFlow sendFileProcessingInfo() { return IntegrationFlows .from(SLACK_CHANNEL) .transform(Message.class, this::prepareSlackMessage) // 仅发送不需要返回,用outboundChannelAdapter更简洁 .handle(WebFlux.outboundChannelAdapter(m -> UriComponentsBuilder.fromUriString(slackConfigurationProperties.getUrl()) .build() .toUri()) .httpMethod(HttpMethod.POST)) .log() .get(); }
如果你需要在多个节点发送Slack通知,只需要在对应节点位置新增.wireTap(SLACK_CHANNEL)即可,无需重复编写发送逻辑。
内容的提问来源于stack exchange,提问作者Luis Pouzo

