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
相关产品推荐
相关产品推荐

