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

