Spring Integration无需轮询读取处理指定路径CSV文件的实现方案
Spring Integration 单文件处理改造方案
核心改动说明
将原有的目录轮询触发逻辑替换为网关触发入口,接收携带file_path(文件路径)、file_name(文件名)消息头的请求,仅读取处理指定的单个文件,后续拆分、聚合、API补全、文件输出逻辑可直接复用原有实现。
具体改造步骤
1. 定义触发网关接口
用于接收处理请求,自动将参数写入消息头:
public interface FileProcessGateway { void processFile(@Header("file_path") String filePath, @Header("file_name") String fileName); }
注册网关Bean到Spring上下文:
@Bean public FileProcessGateway fileProcessGateway(IntegrationFlowContext context) { return context.getGateway(FileProcessGateway.class); }
2. 替换原有轮询入站流定义
将原getUIDsFromTTDandOutputToFile方法的轮询入口替换为网关入口,增加单文件读取校验逻辑:
@Bean @SuppressWarnings("unchecked") public IntegrationFlow getUIDsFromTTDandOutputToFile() { Gson gson = new GsonBuilder().disableHtmlEscaping().create(); return IntegrationFlows // 替换原目录轮询入口,改为网关触发 .from(FileProcessGateway.class) .log(Level.INFO, m -> "TTD UID 2.0 集成处理开始,待处理文件:" + m.getHeaders().get("file_path") + File.separator + m.getHeaders().get("file_name")) // 读取消息头指定的单个文件 .handle((payload, headers) -> new File(headers.get("file_path").toString(), headers.get("file_name").toString())) // 校验文件合法性,可根据需求扩展校验规则 .filter(File::exists, filterSpec -> filterSpec .throwExceptionOnRejection(true) .rejectFlow((message, exception) -> { throw new RuntimeException("指定文件不存在:" + message.getHeaders().get("file_path") + File.separator + message.getHeaders().get("file_name")); }) ) // 后续逻辑与原实现完全一致 .split(Files.splitter()) .channel(c -> c.executor(Executors.newFixedThreadPool(7))) .handle((p, h) -> new CSVUtils().csvColumnSelector((String) p, ttdColNum)) .channel("chunkingChannel") .get(); }
3. 优化聚合分组逻辑适配多文件并发场景
修改CorrelationStrategyIml,用文件名作为聚合分组键,避免同时处理多个文件时数据混乱:
public class CorrelationStrategyIml implements CorrelationStrategy { @Override public Object getCorrelationKey(Message<?> message) { return message.getHeaders().getOrDefault("file_name", "default_group"); } }
4. 删除冗余代码
原轮询场景使用的getFileFilters文件过滤器、轮询配置已无用,可直接删除。
调用方式
需要处理指定文件时,注入FileProcessGateway直接调用即可:
@Autowired private FileProcessGateway fileProcessGateway; // 调用示例 public void triggerFileProcess() { // 传入目标文件的存储路径和文件名 fileProcessGateway.processFile("/your/input/dir", "target.csv"); }
原有
chunker聚合器、enrichFlow补全输出流的代码无需任何修改,输出文件名将自动沿用你传入的file_name加_out.csv后缀,符合原有输出规则。
内容的提问来源于stack exchange,提问作者Jim Dannucci
相关产品推荐
相关产品推荐

