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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 19:25:26