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

Apache Camel拆分后如何保留Header并统计行处理成功/失败指标?

问题

需求为流式读取文件并逐行处理,过程中需跟踪成功处理行数与失败处理行数,且每处理一行后发布对应指标。

但Apache Camel的split组件每次拆分时会创建新Exchange,导致原Exchange的Header丢失,常规聚合机制也不适用(需在拆分过程中实时获取处理结果并发布指标)。原示例路由如下:

onException(Exception.class)
    .handled(true)
    .process(new FailureProcessor()) // 递增failureCount
    .to("{{publishMetrics}}");

from("{{file}}")
    .setHeader("successCount", constant(0))
    .setHeader("failureCount", constant(0))
    .split(body().tokenize("\n"))
    .process(new MyProcessor()) // 处理每行数据
    .process(new SuccessProcessor()) // 递增successCount header
    .to("{{publishMetrics}}") // 发布指标(successCount & failureCount)
    .end()
    .end(); // split结束

解决方案

1. 开启shareUnitOfWork并通过ExchangeProperty存储计数器

split默认不共享原Exchange上下文,开启shareUnitOfWork(true)后,子Exchange会共享父Exchange的UnitOfWork,此时用ExchangeProperty(而非Header)存储计数器,所有子Exchange可共享修改同一个计数器实例:

onException(Exception.class)
    .handled(true)
    .process(exchange -> {
        AtomicInteger failureCount = exchange.getProperty("failureCount", AtomicInteger.class);
        if (failureCount == null) {
            failureCount = new AtomicInteger(0);
            exchange.setProperty("failureCount", failureCount);
        }
        failureCount.incrementAndGet();
    })
    .to("{{publishMetrics}}");

from("{{file}}")
    // 初始化线程安全的原子计数器到ExchangeProperty
    .setProperty("successCount", () -> new AtomicInteger(0))
    .setProperty("failureCount", () -> new AtomicInteger(0))
    // 开启共享上下文
    .split(body().tokenize("\n")).shareUnitOfWork()
        .process(new MyProcessor())
        .process(exchange -> {
            AtomicInteger successCount = exchange.getProperty("successCount", AtomicInteger.class);
            successCount.incrementAndGet();
        })
        .to("{{publishMetrics}}")
    .end();

2. 自定义拆分策略传递共享计数器

若不想依赖UnitOfWork,可自定义Splitter策略,将父Exchange的计数器实例传递给每个子Exchange:

// 自定义拆分策略,负责初始化并传递计数器
class SharedCounterSplitter implements Splitter {
    @Override
    public Iterator<?> split(Exchange exchange, Object body) {
        exchange.setProperty("successCount", new AtomicInteger(0));
        exchange.setProperty("failureCount", new AtomicInteger(0));
        return ((String) body).split("\n").iterator();
    }

    @Override
    public void prepareSubExchange(Exchange exchange, Exchange subExchange, Object value) {
        // 将父Exchange的计数器实例传递给子Exchange
        subExchange.setProperty("successCount", exchange.getProperty("successCount"));
        subExchange.setProperty("failureCount", exchange.getProperty("failureCount"));
    }
}

// 路由配置
onException(Exception.class)
    .handled(true)
    .process(exchange -> {
        AtomicInteger failureCount = exchange.getProperty("failureCount", AtomicInteger.class);
        failureCount.incrementAndGet();
    })
    .to("{{publishMetrics}}");

from("{{file}}")
    .split(new SharedCounterSplitter())
        .process(new MyProcessor())
        .process(exchange -> {
            AtomicInteger successCount = exchange.getProperty("successCount", AtomicInteger.class);
            successCount.incrementAndGet();
        })
        .to("{{publishMetrics}}")
    .end();

3. 使用ThreadLocal全局计数器(单文件场景适用)

针对单文件处理场景,可使用ThreadLocal存储线程安全计数器,避免上下文传递的复杂度:

private static final ThreadLocal<AtomicInteger> SUCCESS_COUNT = ThreadLocal.withInitial(AtomicInteger::new);
private static final ThreadLocal<AtomicInteger> FAILURE_COUNT = ThreadLocal.withInitial(AtomicInteger::new);

onException(Exception.class)
    .handled(true)
    .process(exchange -> FAILURE_COUNT.get().incrementAndGet())
    .to("{{publishMetrics}}");

from("{{file}}")
    .process(exchange -> {
        // 处理文件前重置计数器
        SUCCESS_COUNT.get().set(0);
        FAILURE_COUNT.get().set(0);
    })
    .split(body().tokenize("\n"))
        .process(new MyProcessor())
        .process(exchange -> SUCCESS_COUNT.get().incrementAndGet())
        .to("{{publishMetrics}}")
    .end()
    .process(exchange -> {
        // 文件处理完成后清理ThreadLocal
        SUCCESS_COUNT.remove();
        FAILURE_COUNT.remove();
    });

关键注意事项

  • 禁止用Header存储计数器:split的子Exchange会复制父Header,但子Exchange对Header的修改不会回传,导致计数无法累计。
  • 优先选择shareUnitOfWork方案:这是Camel官方推荐的上下文共享方式,实现简单且符合框架设计规范。
  • 必须用线程安全计数器:若开启split的parallelProcessing并行处理,普通整数会引发线程安全问题,需用AtomicInteger这类原子类。

内容的提问来源于stack exchange,提问作者dark ninja

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 00:53:16