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

如何用Spring Integration轮询器实现目录有文件时仅触发一次Spring Batch

基于Spring Integration轮询器实现单次批量触发Spring Batch的方案

核心思路

默认的文件入站适配器会为每个文件生成单独的消息,会触发多次Batch任务。我们只需要调整为单次轮询拉取当前目录下所有符合条件的文件,再通过聚合器把这批文件整合成单条批量消息,最终只触发一次Batch任务即可。

具体实现步骤

  • 第一步:配置文件入站通道适配器,设置maxMessagesPerPoll = -1,让每轮轮询拉取全部符合条件的文件,同时配置文件过滤器排除未完成写入的文件、目录等不符合要求的内容。
  • 第二步:配置聚合器,基于轮询生成的序列头信息,把同一轮拉取的所有文件消息聚合成单条携带List<File>的消息。
  • 第三步:配置消息处理器,收到聚合后的批量文件消息后,将文件列表作为参数传入Spring Batch任务,仅触发一次作业执行。

代码示例

首先确保项目已引入spring-boot-starter-integration、spring-integration-file、spring-boot-starter-batch相关依赖,以下是核心配置代码:

@Configuration
@EnableIntegration
@EnableBatchProcessing
public class FilePollBatchConfig {

    // 监听的目录路径
    @Value("${polling.dir.path}")
    private File pollingDir;
    // 处理完成后文件归档目录
    @Value("${processed.dir.path}")
    private File processedDir;

    @Autowired
    private JobLauncher jobLauncher;
    @Autowired
    private Job fileProcessJob;

    // 1. 配置文件读取源,添加文件过滤规则
    @Bean
    public MessageSource<File> fileMessageSource() {
        FileReadingMessageSource source = new FileReadingMessageSource();
        source.setDirectory(pollingDir);
        // 组合过滤器:只处理普通文件,且文件最后修改时间距当前大于5s(避免处理未写完的文件)
        source.setFilter(new CompositeFileListFilter<>(Arrays.asList(
                new SimpleFileListFilter(),
                new LastModifiedFileListFilter(5)
        )));
        return source;
    }

    @Bean
    public MessageChannel fileInputChannel() {
        return new DirectChannel();
    }

    // 配置入站适配器,每60秒轮询一次,单次轮询拉取全部符合条件的文件
    @Bean
    @InboundChannelAdapter(channel = "fileInputChannel", poller = @Poller(fixedDelay = "60000", maxMessagesPerPoll = "-1"))
    public MessageSource<File> inboundAdapter() {
        return fileMessageSource();
    }

    // 2. 配置聚合器,合并单次轮询的所有文件为单条消息
    @Bean
    public MessageChannel aggregatedFileChannel() {
        return new DirectChannel();
    }

    @Bean
    @Aggregator(inputChannel = "fileInputChannel", outputChannel = "aggregatedFileChannel",
            releaseStrategy = "releaseStrategy", correlationStrategy = "correlationStrategy")
    public List<File> aggregateFiles(List<Message<File>> messages) {
        return messages.stream()
                .map(Message::getPayload)
                .collect(Collectors.toList());
    }

    @Bean
    public ReleaseStrategy releaseStrategy() {
        // 收到的消息数等于本次轮询拉取的总文件数时释放聚合结果
        return new SimpleSequenceSizeReleaseStrategy();
    }

    @Bean
    public CorrelationStrategy correlationStrategy() {
        // 同一轮轮询的消息归属同一个聚合组
        return message -> message.getHeaders().get(IntegrationMessageHeaderAccessor.CORRELATION_ID);
    }

    // 3. 配置消息处理器,触发单次Spring Batch任务
    @Bean
    @ServiceActivator(inputChannel = "aggregatedFileChannel")
    public MessageHandler batchJobLauncherHandler() {
        return message -> {
            List<File> filesToProcess = (List<File>) message.getPayload();
            // 构造Job参数,传入文件列表
            JobParameters jobParameters = new JobParametersBuilder()
                    .addLong("startTime", System.currentTimeMillis())
                    .addJobParameter("files", filesToProcess, List.class)
                    .toJobParameters();
            // 仅触发一次Batch任务
            jobLauncher.run(fileProcessJob, jobParameters);
            // 处理完成后归档文件,也可配置source的deleteSourceFiles属性直接删除原文件
            filesToProcess.forEach(file -> file.renameTo(new File(processedDir, file.getName())));
        };
    }
}

注意事项

  • 可根据业务需求调整轮询规则,支持fixedDelay固定间隔、cron表达式两种配置方式。
  • 已处理的文件建议归档或删除,避免下次轮询重复处理。
  • 若需要异步执行Batch任务,可自定义异步JobLauncher,避免阻塞轮询线程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 13:24:04