Apache Camel流式拆分(tokenize)场景下如何识别并处理最后元素?
解决方案:Apache Camel流式处理中识别流结束
方法一:使用Splitter的onCompletion回调(推荐)
end()方法仅用于闭合Split的子路由块,无法触发流处理完成后的逻辑。Camel的Splitter组件提供了onCompletion()回调,会在所有拆分元素处理完成后触发,完美适配流式场景,且支持所有可被Splitter处理的数据源(文件、HTTP请求体、S3/SFTP对象等),不会导致内存占用激增。
完整路由示例:
from("file://target/input/?delete=true") .log("Started processing [${header.CamelFileNameOnly}] ...") .split(body().tokenize("\n")).streaming() .log("ELEMENT ${body}") .process(this::doStuffWithElement) // 流处理完成后的收尾逻辑 .onCompletion() .process(exchange -> { String fileName = exchange.getIn().getHeader("CamelFileNameOnly", String.class); // 执行后续操作:比如发送通知、归档文件、更新状态等 exchange.getIn().setBody("All lines processed for file: " + fileName); }) .end(); // 闭合onCompletion块
方法二:跟踪最后一个元素(处理最后一行时立即执行逻辑)
如果需要在处理最后一行的同时触发操作,可以利用Camel内置的Exchange.SPLIT_COMPLETE属性,结合自定义AggregationStrategy标记最后一个元素:
- 定义跟踪策略类:
public class LastElementTracker implements AggregationStrategy { @Override public Exchange aggregate(Exchange oldExchange, Exchange newExchange) { // 检查当前元素是否为最后一个 Boolean isLast = newExchange.getProperty(Exchange.SPLIT_COMPLETE, Boolean.class); if (Boolean.TRUE.equals(isLast)) { newExchange.setProperty("IsLastLine", true); } return newExchange; } }
- 修改路由:
from("file://target/input/?delete=true") .log("Started processing [${header.CamelFileNameOnly}] ...") // 传入自定义AggregationStrategy跟踪最后一行 .split(body().tokenize("\n"), new LastElementTracker()).streaming() .log("ELEMENT ${body}") .process(exchange -> { doStuffWithElement(exchange); // 判断是否为最后一行并执行逻辑 Boolean isLast = exchange.getProperty("IsLastLine", Boolean.class); if (Boolean.TRUE.equals(isLast)) { log.info("Processing final line: ${body}"); // 执行最后一行专属操作 } }) .end();
关键说明
- 两种方案均为流式处理,仅加载当前行到内存,避免内存占用过高;
- 适配所有支持Camel Splitter的数据源,无需针对文件、HTTP、S3/SFTP做特殊适配;
Exchange.SPLIT_COMPLETE是Camel内置属性,拆分到最后一个元素时会自动设为true,可靠性有保障。
内容的提问来源于stack exchange,提问作者Guido
相关产品推荐
相关产品推荐

