Spring Integration升级后scatterGather的aggregator返回结果失败如何解决
Spring Integration升级后aggregator返回结果异常解决方案
问题根因
Spring Integration从旧版升级到5.x版本后,aggregator(scatterGather中的gatherer)的输出逻辑发生了变更:默认情况下如果outputProcessor返回的是Collection类型实例,框架会默认判定你需要将集合中的每个元素作为独立的Message分别输出,因此会校验集合内每个元素都是合法的Message类型。你返回的List<String>属于集合类型,触发了这个校验逻辑,因此抛出了类型不匹配异常。
解决方案
有两种方案可选,任选其一即可恢复1.2.x版本的原有行为:
方案1(推荐:无需修改原有业务逻辑)
在gatherer配置中添加outputIsSingleResult(true)配置,强制框架将返回的集合整体作为单个Message的payload处理,不做拆分校验:
public IntegrationFlow stringTransformer() { return f -> f .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(f1 -> f1.<List<String>, List<String>>transform(p -> p.stream().map(String::toUpperCase) .collect(Collectors.toList()))) .recipientFlow(f2 -> f2.<List<String>, List<String>>transform(p -> p.stream().map(String::toLowerCase) .collect(Collectors.toList()))), gatherer -> gatherer.expireGroupsUponCompletion(true) .outputIsSingleResult(true) // 新增这一行配置 .outputProcessor(messageGroup -> { List<String> finalList; finalList = messageGroup.getMessages().stream() .map(Message::getPayload) .map(p -> (List<String>)p) .flatMap(Collection::stream) .collect(Collectors.toList()); return finalList; }), s -> s.errorChannel("scatterGatherErrorChannel")); }
方案2:手动将返回结果包装为Message对象
修改outputProcessor的返回值,将最终生成的List包装为Message实例:
.outputProcessor(messageGroup -> { List<String> finalList = messageGroup.getMessages().stream() .map(Message::getPayload) .map(p -> (List<String>)p) .flatMap(Collection::stream) .collect(Collectors.toList()); // 手动包装为Message对象 return MessageBuilder.withPayload(finalList).build(); })
内容的提问来源于stack exchange,提问作者kiran reddy
相关产品推荐
相关产品推荐

