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

如何在Spring Integration FTP中聚合单次轮询文件并压缩发邮件?

Spring Integration FTP 单次轮询文件聚合压缩后发邮件实现方案

方案一:利用Aggregator聚合单次轮询消息

Spring Integration的Aggregator组件可直接实现单次轮询消息的聚合,核心是通过轮询关联ID识别同一次轮询的所有文件消息,当轮询完成时触发聚合操作。

关键修改步骤

  1. 替换原有的publishSubscribeChannel为aggregate组件,配置聚合规则
  2. 实现聚合后的文件压缩逻辑
  3. 调整邮件发送逻辑,改为发送压缩后的ZIP文件

修改后的代码示例

public IntegrationFlow createIntFlow(MessageDirectory directory, ConcurrentMetadataStore metadataStore) {
    var directoryName = directory.getDirectoryName();

    return IntegrationFlows
            .from(getInboundAdapter(directory, metadataStore),
                    e -> e.id(directoryName + "-PerPoller")
                            .autoStartup(true)
                            .poller(Pollers
                                    .cron(directory.getSchedule())
                                    .taskExecutor(simpleAsyncTaskExecutor)
                                    .errorChannel("errorChannel")
                                    .maxMessagesPerPoll(-1)))
            .log()
            // 按单次轮询的关联ID分组,轮询结束时触发聚合
            .aggregate(a -> a
                    .correlationStrategy(m -> m.getHeaders().get(IntegrationMessageHeaderAccessor.POLLER_METADATA))
                    .releaseStrategy(g -> g.isComplete())
                    .sendPartialResultOnExpiry(false)
                    .outputProcessor(this::compressFilesToZip))
            // 发送包含ZIP文件的邮件
            .handle(c -> emailEmitter.createZipFileEvent(directory, c.getPayload().toString()))
            .get();
}

// 聚合结果处理器:将多个文件压缩为ZIP
private File compressFilesToZip(MessageGroup messageGroup) {
    List<File> files = messageGroup.getMessages().stream()
            .map(m -> (File) m.getPayload())
            .collect(Collectors.toList());
    
    // 创建临时ZIP文件并执行压缩
    File zipFile = new File(System.getProperty("java.io.tmpdir") + "/" + directory.getDirectoryName() + "_" + System.currentTimeMillis() + ".zip");
    try (ZipOutputStream zos = new ZipOutputStream(new FileOutputStream(zipFile))) {
        for (File file : files) {
            try (FileInputStream fis = new FileInputStream(file)) {
                ZipEntry zipEntry = new ZipEntry(file.getName());
                zos.putNextEntry(zipEntry);
                byte[] buffer = new byte[1024];
                int length;
                while ((length = fis.read(buffer)) > 0) {
                    zos.write(buffer, 0, length);
                }
                zos.closeEntry();
            }
        }
    } catch (IOException e) {
        throw new RuntimeException("压缩文件失败", e);
    }
    return zipFile;
}

核心说明

  • correlationStrategy:使用POLLER_METADATA头作为关联ID,确保同一次轮询的消息被分到同一组
  • releaseStrategy:g.isComplete()表示当轮询的所有消息都被接收后,触发聚合
  • outputProcessor:自定义处理逻辑,将聚合后的文件列表压缩为ZIP文件

方案二:临时目录暂存+轮询结束事件触发

如果更倾向于先将文件移到临时目录再统一处理,可通过监听轮询事件实现。

关键修改步骤

  1. 添加文件移动处理器,将单次轮询的文件移到指定临时目录
  2. 监听轮询结束事件,触发压缩操作
  3. 压缩完成后发送邮件并清理临时目录

修改后的代码示例

public IntegrationFlow createIntFlow(MessageDirectory directory, ConcurrentMetadataStore metadataStore, ApplicationEventPublisher publisher) {
    var directoryName = directory.getDirectoryName();
    File tempDir = new File(System.getProperty("java.io.tmpdir") + "/" + directoryName + "_temp");
    if (!tempDir.exists()) {
        tempDir.mkdirs();
    }

    // 监听轮询开始/结束事件
    publisher.addApplicationListener((ApplicationEvent event) -> {
        if (event instanceof PollerPollingStartedEvent) {
            // 轮询开始前清空临时目录
            Arrays.stream(tempDir.listFiles()).forEach(File::delete);
        } else if (event instanceof PollerPollingCompletedEvent) {
            PollerPollingCompletedEvent completedEvent = (PollerPollingCompletedEvent) event;
            if (completedEvent.getPollerMetadata().getId().equals(directoryName + "-PerPoller")) {
                // 轮询结束,压缩临时目录文件
                File zipFile = compressTempDirToZip(tempDir, directoryName);
                // 发送邮件
                emailEmitter.createZipFileEvent(directory, zipFile.getAbsolutePath());
                // 清理临时文件
                zipFile.delete();
            }
        }
    });

    return IntegrationFlows
            .from(getInboundAdapter(directory, metadataStore),
                    e -> e.id(directoryName + "-PerPoller")
                            .autoStartup(true)
                            .poller(Pollers
                                    .cron(directory.getSchedule())
                                    .taskExecutor(simpleAsyncTaskExecutor)
                                    .errorChannel("errorChannel")
                                    .maxMessagesPerPoll(-1)))
            .log()
            // 将文件移动到临时目录
            .handle(Files.outboundGateway(tempDir)
                    .fileExistsMode(FileExistsMode.REPLACE)
                    .autoCreateDirectory(true))
            .get();
}

// 压缩临时目录为ZIP
private File compressTempDirToZip(File tempDir, String directoryName) {
    File zipFile = new File(System.getProperty("java.io.tmpdir") + "/" + directoryName + "_" + System.currentTimeMillis() + ".zip");
    try (ZipOutputStream zos = new ZipOutputStream(new FileOutputStream(zipFile))) {
        for (File file : Objects.requireNonNull(tempDir.listFiles())) {
            try (FileInputStream fis = new FileInputStream(file)) {
                ZipEntry zipEntry = new ZipEntry(file.getName());
                zos.putNextEntry(zipEntry);
                byte[] buffer = new byte[1024];
                int length;
                while ((length = fis.read(buffer)) > 0) {
                    zos.write(buffer, 0, length);
                }
                zos.closeEntry();
                file.delete();
            }
        }
    } catch (IOException e) {
        throw new RuntimeException("压缩临时目录失败", e);
    }
    return zipFile;
}

核心说明

  • 利用Spring Integration的Files.outboundGateway实现文件移动
  • 通过监听PollerPollingStartedEvent和PollerPollingCompletedEvent识别轮询的开始与结束
  • 轮询开始时清空临时目录,结束时压缩目录内文件并发送邮件

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 07:13:20