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

Apache Camel Splitter处理CSV分块、表头复用及groupBy配置咨询

优化CSV拆分逻辑与实现groupBy参数可配置

咱们来一步步解决你遇到的两个核心问题:让处理器代码更优雅安全,同时实现groupBy参数从Camel上下文动态读取。

一、优雅化处理器代码(解决线程安全与逻辑冗余)

你当前的处理器用了全局变量hdrFlag、HEADER、cntr,这在多线程处理文件时会引发并发冲突(比如多个文件的状态互相干扰)。我们可以把这些状态转移到Exchange属性中,同时简化文本处理逻辑:

修改后的处理器代码

import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.StringReader;
import java.io.BufferedReader;

public class AsciiSplitterProcessor implements Processor {
    private static final Logger log = LoggerFactory.getLogger(AsciiSplitterProcessor.class);
    // 定义常量避免硬编码属性键
    private static final String CSV_HEADER = "CSV_HEADER";
    private static final String SPLIT_COUNTER = "SPLIT_COUNTER";

    @Override
    public void process(Exchange exchange) throws Exception {
        log.info("Ascii Splitter Processor :: start");
        String inBody = exchange.getIn().getBody(String.class);
        String fileName = exchange.getIn().getHeader("CamelFileName", String.class);
        
        // 初始化计数器(仅当前文件首次处理时)
        if (exchange.getProperty(SPLIT_COUNTER) == null) {
            exchange.setProperty(SPLIT_COUNTER, 1);
        }
        int cntr = exchange.getProperty(SPLIT_COUNTER, Integer.class);
        
        // 生成规范的拆分文件名
        String filePrefix = fileName.substring(0, fileName.lastIndexOf("."));
        String fileSuffix = fileName.substring(fileName.lastIndexOf("."));
        String newFileName = String.format("%s_%d%s", filePrefix, cntr, fileSuffix);
        exchange.getIn().setHeader("CamelFileName", newFileName);
        log.info("File being processed: {}", newFileName);
        
        // 处理表头与内容拼接
        StringBuilder sb = new StringBuilder();
        try (BufferedReader reader = new BufferedReader(new StringReader(inBody))) {
            String firstLine = reader.readLine();
            // 首次处理:提取表头并存储到Exchange属性
            if (exchange.getProperty(CSV_HEADER) == null) {
                exchange.setProperty(CSV_HEADER, firstLine + "\r\n");
                // 第一块内容保留原始行(包含表头)
                sb.append(firstLine).append("\r\n");
                // 读取剩余行
                String line;
                while ((line = reader.readLine()) != null) {
                    sb.append(line).append("\r\n");
                }
            } else {
                // 后续块:追加表头+当前内容
                sb.append(exchange.getProperty(CSV_HEADER, String.class));
                sb.append(firstLine).append("\r\n");
                String line;
                while ((line = reader.readLine()) != null) {
                    sb.append(line).append("\r\n");
                }
            }
        }
        
        exchange.getIn().setBody(sb.toString());
        // 计数器自增
        exchange.setProperty(SPLIT_COUNTER, cntr + 1);
        log.debug("Processed content: {}", sb.toString());
    }
}

关键优化点:

  • 线程安全:把表头、计数器从全局变量移到Exchange属性,每个文件的处理状态完全独立,避免并发冲突。
  • 资源安全:用try-with-resources自动关闭BufferedReader,杜绝资源泄漏。
  • 逻辑清晰:拆分首次提取表头和后续追加表头的逻辑,代码可读性大幅提升。
  • 规范命名:使用格式化字符串生成文件名,定义常量代替硬编码的属性键。

二、实现groupBy参数的可配置化

要让tokenize的group参数从Camel上下文读取,有两种实用方式:

方式1:配置文件占位符(推荐)

在你的Camel配置文件(比如camel.properties或application.properties)中添加:

split.group.size=50

然后修改路由中的tokenize节点,用占位符引用配置值:

<split streaming="true">
    <tokenize token="\n" group="{{split.group.size}}"/>
    <!-- 后续逻辑不变 -->
</split>

方式2:动态读取上下文属性(支持运行时调整)

如果需要在运行时动态修改分组大小,可以用simple表达式读取Camel上下文属性,同时设置默认值:

<split streaming="true">
    <tokenize token="\n" group="${camelContext.properties.getProperty('split.group.size', '50')}"/>
    <!-- 后续逻辑不变 -->
</split>

如果上下文属性不存在,会自动使用默认值50。

修改后的完整路由

<!-- route on Weekends -->
<route id="inRouteWkEndBfmt1" routePolicyRef="startPolicyWkEnd" autoStartup="false" >
    <from id="mainProcessingRouteWkEnd" ref="AsciiGatewayBackfillmt1" />
    <convertBodyTo type="java.lang.String" />
    <log message="File ${file:name} was received."/>
    <setHeader headerName="messageDateTime">
        <simple>${date:now:MM-dd-yyyy-HH:mm:ss}</simple>
    </setHeader>
    <!-- 使用可配置的分组参数 -->
    <split streaming="true" >
        <tokenize token="\n" group="{{split.group.size}}"/>
        <log message="Split line Body: ${body}"/>
        <process ref="asciiSplitterProcessor" />
        <log loggingLevel="INFO" message="Successfully sent ${file:name} to MT1 Core for Analytics Observation." />
        <to id="windowsShareTargetWkEnd" uri="file://{{target.folder}}" />
    </split>
    <process ref="asciiCleanUp" />
</route>

额外小建议

  • 可以给AsciiSplitterProcessor添加一个defaultHeader属性,作为表头提取失败时的 fallback,增强鲁棒性。
  • 你已经开启了streaming="true",这个配置能避免大文件内存溢出,保持即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:18:04