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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 23:06:08