You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 18:35:53