Spring Integration拆分聚合流异常处理:终止当前链并继续处理元素
需求描述
希望实现一个包含split/aggregate的Spring Integration集成流,要求:拆分后的元素在处理器中抛出异常时,终止该元素后续的处理器执行,但继续处理其他拆分元素。
初始尝试代码
IntegrationFlow.from( WebFlux.inboundGateway("/jira/version") .requestMapping(r -> r.methods(HttpMethod.POST) .consumes("application/json")) .requestPayloadType(String.class) .replyChannel(replyChannel) .errorChannel(errorChannel) .mappedRequestHeaders(parameter.getJiraHeaderSignature())) .split() .handle(handler1, s -> s.advice(advice())) .handle(handler2, s -> s.advice(advice())) .handle(handler3, s -> s.advice(advice())) .aggregate() .get(); public Advice advice() { var advice = new ExpressionEvaluatingRequestHandlerAdvice(); advice.setReturnFailureExpressionResult(true); advice.setOnFailureExpressionString("payload"); return advice; }
场景说明
处理包含2个元素的集合时:
- 第一个元素在handler1抛出异常,要求handler2、handler3不再执行;
- 第二个元素正常执行所有处理器;
- 聚合时需返回包含2个原始消息的集合。
但启用TrapException时,虽能终止当前处理器链,但聚合无法正常结束。
更新1:使用Detour解决链终止问题
通过ExpressionEvaluatingRequestHandlerAdvice配置failureChannel,可实现异常元素终止后续处理器,但暂未找到在聚合器内直接过滤ErrorMessage的方法,只能在聚合后通过handle过滤:
.handle(handler1, s -> s.advice(advice())) .handle(handler2, s -> s.advice(advice())) .channel("jiraAggregatorInputChannel") .aggregate() .<List>handle((p, h) -> p.stream() .filter(i -> i instanceof CustomObject) .toList()) .get() public Advice advice() { var advice = new ExpressionEvaluatingRequestHandlerAdvice(); advice.setTrapException(true); advice.setFailureChannelName("jiraAggregatorInputChannel"); return advice; }
更新2:聚合后过滤与聚合器内过滤的行为差异
聚合后过滤代码
IntegrationFlow.from( WebFlux.inboundGateway("/jira/version") .requestMapping(r -> r.methods(HttpMethod.POST) .consumes("application/json")) .requestPayloadType(String.class) .replyChannel(replyChannel) .errorChannel(errorChannel)) .handle(handler1, s -> s.advice(advice())) .handle(handler2, s -> s.advice(advice())) .channel("jiraAggregatorInputChannel") .aggregate() .<List>handle((p, h) -> p.stream() .filter(i -> i instanceof CustomObject) .toList()) .get()
结果:返回纯CustomObject列表
[{"issuesNumberAkuiteo": ...., "idVersionJira": ....}]
聚合器内过滤代码
IntegrationFlow.from( WebFlux.inboundGateway("/jira/version") .requestMapping(r -> r.methods(HttpMethod.POST) .consumes("application/json")) .requestPayloadType(String.class) .replyChannel(replyChannel) .errorChannel(errorChannel)) .handle(handler1, s -> s.advice(advice())) .handle(handler2, s -> s.advice(advice())) .channel("jiraAggregatorInputChannel") .aggregate(a -> a.outputProcessor(group -> group .getMessages() .stream() .filter(i -> i.getPayload() instanceof CustomObject) .toList())) .get()
结果:返回GenericMessage列表(包含payload和headers),即使WebFlux默认extractReplyPayload=true,行为仍不一致
[{"payload": { .... }, "headers": { .... }}]
更新3:聚合器内映射Payload的异常
尝试在聚合器内先映射Payload再过滤,抛出如下异常:
IntegrationFlow.from( WebFlux.inboundGateway("/jira/version") .requestMapping(r -> r.methods(HttpMethod.POST) .consumes("application/json")) .requestPayloadType(String.class) .replyChannel(replyChannel) .errorChannel(errorChannel)) .handle(handler1, s -> s.advice(advice())) .handle(handler2, s -> s.advice(advice())) .channel("jiraAggregatorInputChannel") .aggregate(a -> a.outputProcessor(group -> group .getMessages() .stream() .map(Message::getPayload) .filter(i -> i instanceof CustomObject) .toList())) .get()
异常信息:
The expected collection of Messages contains non-Message element: class fr.eksae.erp.connector.model.CustomObject: class fr.eksae.erp.connector.model.CustomObject
最终解决方案
通过在聚合器的outputProcessor中手动创建GenericMessage,包装过滤并映射后的Payload列表,解决返回格式问题:
IntegrationFlow.from( WebFlux.inboundGateway("/jira/version") .requestMapping(r -> r.methods(HttpMethod.POST) .consumes("application/json")) .requestPayloadType(String.class) .replyChannel(replyChannel) .errorChannel(errorChannel)) .split() .handle(handler1, s -> s.advice(advice())) .handle(handler2, s -> s.advice(advice())) .channel("jiraAggregatorInputChannel") .aggregate(a -> a.outputProcessor( group -> new GenericMessage<>(group .getMessages() .stream() .filter(i -> i.getPayload() instanceof CustomObject) .map(Message::getPayload) .toList()))) .get() public Advice advice() { var advice = new ExpressionEvaluatingRequestHandlerAdvice(); advice.setTrapException(true); advice.setFailureChannelName("jiraAggregatorInputChannel"); return advice; }
内容的提问来源于stack exchange,提问作者Médéric Martin
相关产品推荐
相关产品推荐

