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

如何在Apache Camel中按自定义条件停止Splitter拆分流程

Apache Camel Splitter 自定义条件终止拆分流程方案

Splitter自带的无参stopOnException()、stopOnAggregateException()仅支持全量异常触发停止,无法实现差异化异常处理、自定义业务规则终止,以下两种方案可直接满足需求。

之前尝试的stop()方法仅能终止当前分片的子路由执行,作用域不覆盖Splitter的迭代逻辑,因此无法终止整个拆分流程。

方案1:自定义聚合策略(兼容所有Camel版本,支持任意自定义终止规则)

该方案通过Splitter的聚合回调逻辑判断终止条件,手动设置Camel内置的拆分完成标记,不仅支持按异常类型判断,还支持按消息内容、业务规则等非异常场景触发终止,兼容性最强。
实现步骤:

  • 配置路由级差异化异常捕获:对LineProcessingException设置continued(true),直接跳过当前行继续后续处理;对IOException设置专属终止标记,标记异常已处理避免向上抛出
  • 实现自定义AggregationStrategy,每次聚合分片结果时检查终止标记,若标记存在则设置Exchange.SPLIT_COMPLETE为true,Splitter识别到该标记后会立刻停止遍历剩余分片
  • 拆分时绑定自定义聚合策略,不要调用无参stopOnException()

完整实现代码:

import org.apache.camel.AggregationStrategy;
import org.apache.camel.Exchange;
import org.apache.camel.LoggingLevel;
import java.util.ArrayList;
import java.util.List;

// 自定义拆分聚合策略
AggregationStrategy splitControlStrategy = new AggregationStrategy() {
    @Override
    public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
        // 首个分片处理时oldExchange为空,直接返回当前Exchange
        if (oldExchange == null) {
            return newExchange;
        }
        // 检查是否触发终止标记
        Boolean needStop = oldExchange.getProperty("STOP_SPLIT", false, Boolean.class);
        if (Boolean.TRUE.equals(needStop)) {
            // 设置Splitter内置终止标识,立刻停止剩余分片处理
            newExchange.setProperty(Exchange.SPLIT_COMPLETE, true);
        }
        // 以下为正常结果聚合逻辑,可根据自身业务调整
        Object lineResult = newExchange.getIn().getBody();
        if (lineResult != null) {
            List<Object> totalResults = oldExchange.getIn().getBody(List.class);
            if (totalResults == null) {
                totalResults = new ArrayList<>();
            }
            totalResults.add(lineResult);
            oldExchange.getIn().setBody(totalResults);
        }
        return oldExchange;
    }
};

// 业务路由
from("timer://file-poll?period=5s")
    // 行处理异常:跳过当前行继续
    .onException(LineProcessingException.class)
        .log(LoggingLevel.WARN, "跳过异常数据行: ${exception.message}")
        .continued(true)
    .end()
    // IO异常:标记终止拆分
    .onException(IOException.class)
        .log(LoggingLevel.ERROR, "IO异常,终止剩余行处理: ${exception.message}")
        .setProperty("STOP_SPLIT", constant(true))
        .handled(true)
    .end()
    .to("direct:load-file")
    // 按换行拆分,绑定自定义聚合策略
    .split(body().tokenize("\n"), splitControlStrategy)
        .to("direct:process-line")
    .end()
    .to("direct:save-result");

方案2:指定异常类型触发停止(Camel 2.16+版本适用,实现更轻量)

Camel 2.16及以上版本的stopOnException支持传入指定异常类型,仅当抛出匹配类型的异常时才会终止拆分,其他异常不会触发停止,实现代码量更少。
实现步骤:

  • 定义一个无业务逻辑的标识性异常StopSplitException,专门用来触发拆分终止
  • 配置差异化异常捕获:LineProcessingException设置continued(true)跳过;IOException捕获后包装为StopSplitException抛出
  • Splitter配置.stopOnException(StopSplitException.class),仅响应指定终止异常

完整实现代码:

// 自定义拆分终止标识异常,无需额外逻辑
public class StopSplitException extends RuntimeException {
    public StopSplitException(String message, Throwable cause) {
        super(message, cause);
    }
}

// 业务路由
from("timer://file-poll?period=5s")
    // 行处理异常:跳过当前行继续
    .onException(LineProcessingException.class)
        .log(LoggingLevel.WARN, "跳过异常数据行: ${exception.message}")
        .continued(true)
    .end()
    // IO异常:包装为终止异常抛出
    .onException(IOException.class)
        .log(LoggingLevel.ERROR, "IO异常,终止剩余行处理: ${exception.message}")
        .throwException(StopSplitException.class, "拆分流程强制终止", "${exception}")
    .end()
    .to("direct:load-file")
    // 仅当抛出StopSplitException时停止拆分
    .split(body().tokenize("\n")).stopOnException(StopSplitException.class)
        .to("direct:process-line")
    .end()
    .to("direct:save-result");

选型建议:

  • 如果需要基于消息内容、自定义业务规则(非异常场景)触发拆分终止,或使用的Camel版本低于2.16,选择方案1
  • 如果仅需要按异常类型差异化控制,且Camel版本满足要求,选择方案2,维护成本更低

内容的提问来源于stack exchange,提问作者Eugene M

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 21:06:23