基于Apache Camel解析CSV并统计行数及处理错误行
Apache Camel原生实现CSV解析统计与错误行记录
完全可以通过Apache Camel的原生组件和特性实现需求,核心用到Bindy行级错误处理、Split聚合策略、Exchange属性统计以及File组件写入错误文件,无需额外第三方依赖。
实现步骤与核心代码
1. 配置Bindy CSV格式的行级错误处理器
Bindy原生支持捕获单条CSV行的解析错误,通过自定义BindyErrorHandlerSupport可以直接获取错误行内容、行号和异常信息,同时统计错误数:
BindyCsvDataFormat bindyCsv = new BindyCsvDataFormat(CustomObject.class); bindyCsv.setErrorHandler(new BindyErrorHandlerSupport() { @Override public void onError(Exception exception, String record, int rowNumber) throws Exception { Exchange exchange = getCurrentExchange(); // 初始化或获取错误行列表 List<String> errorRows = exchange.getProperty("errorRows", List.class); if (errorRows == null) { errorRows = new ArrayList<>(); exchange.setProperty("errorRows", errorRows); } // 存入带行号的错误行(格式可自定义) errorRows.add(String.format("%d,%s", rowNumber, record)); // 累加错误行数 int currentErrorCount = exchange.getProperty("errorCount", Integer.class); exchange.setProperty("errorCount", currentErrorCount + 1); } });
2. 自定义聚合策略合并批次统计数据
因为你用了按100行拆分的批量处理,需要通过AggregationStrategy把每个批次的成功/错误统计、错误行合并到父Exchange中:
public class CsvStatsAggregationStrategy implements AggregationStrategy { @Override public Exchange aggregate(Exchange oldExchange, Exchange newExchange) { if (oldExchange == null) { // 第一个批次,初始化全局统计属性 newExchange.setProperty("successCount", newExchange.getProperty("successCount", Integer.class)); newExchange.setProperty("errorCount", newExchange.getProperty("errorCount", Integer.class)); newExchange.setProperty("errorRows", newExchange.getProperty("errorRows", List.class)); return newExchange; } // 累加成功行数 int totalSuccess = oldExchange.getProperty("successCount", Integer.class) + newExchange.getProperty("successCount", Integer.class); oldExchange.setProperty("successCount", totalSuccess); // 累加错误行数 int totalError = oldExchange.getProperty("errorCount", Integer.class) + newExchange.getProperty("errorCount", Integer.class); oldExchange.setProperty("errorCount", totalError); // 合并错误行列表 List<String> allErrors = oldExchange.getProperty("errorRows", List.class); allErrors.addAll(newExchange.getProperty("errorRows", List.class)); oldExchange.setProperty("errorRows", allErrors); return oldExchange; } }
3. 完整路由配置
整合上述组件,实现从文件读取、批量解析、统计、错误行写入的完整流程:
from(inputFileUri) .routeId(CUSTOM_ROUTE_ID) // 初始化全局统计属性 .setProperty("successCount", constant(0)) .setProperty("errorCount", constant(0)) .setProperty("errorRows", method(ArrayList.class, "new")) // 按100行拆分CSV,使用自定义聚合策略合并统计数据 .split(body().tokenize("\n", 100, true), new CsvStatsAggregationStrategy()) .shareUnitOfWork(true) // 反序列化CSV到CustomObject列表 .unmarshal(bindyCsv) .convertBodyTo(List.class) // 统计当前批次成功解析行数 .process(exchange -> { List<CustomObject> validItems = exchange.getIn().getBody(List.class); int batchSuccess = validItems.size(); exchange.setProperty("successCount", batchSuccess); }) // 执行自定义业务处理 .process(customProcessor) .end() // 结束拆分逻辑 // 处理完成后,写入错误文件(如果存在错误行) .choice() .when(simple("${property.errorCount} > 0")) .setBody(simple("${property.errorRows}")) .marshal().csv() // 将错误行列表序列化为CSV格式 .to("file:///your/error/dir?fileName=csv-errors-${date:now:yyyyMMddHHmmss}.csv") .end() // 打印最终统计结果 .log("CSV处理完成:成功解析${property.successCount}行,错误${property.errorCount}行");
4. 补充:全局异常捕获
针对Bindy未捕获的致命异常(如IO错误、格式完全非法的批次),可以添加全局异常处理:
onException(Exception.class) .process(exchange -> { String rawBatch = exchange.getIn().getBody(String.class); List<String> errorRows = exchange.getProperty("errorRows", List.class); if (errorRows == null) { errorRows = new ArrayList<>(); exchange.setProperty("errorRows", errorRows); } // 记录批次异常信息 errorRows.add(String.format("批次异常:%s\n原始内容:%s", exchange.getException().getMessage(), rawBatch)); // 统计该批次所有行错误 int batchSize = rawBatch.split("\n").length; exchange.setProperty("errorCount", exchange.getProperty("errorCount", Integer.class) + batchSize); }) .handled(true);
关键特性说明
- Bindy行级错误处理:精准捕获单条CSV行的解析错误,不会因为某一行失败导致整个批次丢弃
- Split聚合策略:确保批量处理时,所有批次的统计数据和错误行能合并到全局Exchange中
- Exchange属性:原生的属性机制实现统计计数,无需额外存储组件
- File组件动态命名:通过
${date:now}生成唯一错误文件名,避免覆盖
内容的提问来源于stack exchange,提问作者Geeky
相关产品推荐
相关产品推荐

