如何在Spring Integration FTP中聚合单次轮询文件并压缩发邮件?
Spring Integration FTP 单次轮询文件聚合压缩后发邮件实现方案
方案一:利用Aggregator聚合单次轮询消息
Spring Integration的Aggregator组件可直接实现单次轮询消息的聚合,核心是通过轮询关联ID识别同一次轮询的所有文件消息,当轮询完成时触发聚合操作。
关键修改步骤
- 替换原有的
publishSubscribeChannel为aggregate组件,配置聚合规则 - 实现聚合后的文件压缩逻辑
- 调整邮件发送逻辑,改为发送压缩后的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文件
方案二:临时目录暂存+轮询结束事件触发
如果更倾向于先将文件移到临时目录再统一处理,可通过监听轮询事件实现。
关键修改步骤
- 添加文件移动处理器,将单次轮询的文件移到指定临时目录
- 监听轮询结束事件,触发压缩操作
- 压缩完成后发送邮件并清理临时目录
修改后的代码示例
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
相关产品推荐
相关产品推荐

