Spring Integration SFTP分批次流式上传文件时无法追加覆盖问题及PipedInput/OutputStream使用咨询
首先,咱们先明确核心问题:你当前用SftpRemoteFileTemplate.send()方法上传时,每次都会覆盖目标文件,这是因为该方法默认执行的是SFTP的PUT操作——要么创建新文件,要么直接覆盖已存在的文件,完全不支持追加模式。另外你用的PipedInputStream/PipedOutputStream是一次性的线程间通信工具,每次调用write()方法都会新建一套流,上传的自然是独立的小批量数据,最终覆盖之前的文件。
下面给你两种针对性的解决方案,按需选择:
方案一:小批量数据直接追加到SFTP文件(推荐)
这种方式不需要复杂的流维护,每次处理完一批数据后,直接通过SFTP的追加模式写入目标文件。核心是绕过send()方法,直接调用底层ChannelSftp的API:
public void write(List<? extends Item> items) throws Exception { log.debug("Write {}", items); try (ByteArrayOutputStream baos = new ByteArrayOutputStream()) { // 先把当前批次数据序列化为字节数组 csvSerializer.serialize(baos, items.stream()); byte[] batchData = baos.toByteArray(); // 获取SFTP会话,执行追加操作 sftpRemoteFileTemplate.execute(session -> { ChannelSftp channel = session.getClientInstance(); String remoteFilePath = contactSenderSftpProperties.getSftpSessionProperties().getBaseSftpPath() + "/contacts.csv"; try { // 检查目标文件是否存在 channel.lstat(remoteFilePath); // 文件存在,用APPEND模式追加数据 channel.put(new ByteArrayInputStream(batchData), remoteFilePath, ChannelSftp.APPEND); } catch (SftpException e) { if (e.id == ChannelSftp.SSH_FX_NO_SUCH_FILE) { // 文件不存在,直接创建新文件写入 channel.put(new ByteArrayInputStream(batchData), remoteFilePath); } else { throw new RuntimeException("SFTP操作失败", e); } } return null; }); } }
这种方式的好处是简单可靠,不需要维护全局流,每批数据处理完直接提交到SFTP,也不用纠结Piped流的问题。
方案二:大文件流式持续写入(适合超大数据量)
如果你的数据量极大,无法一次性放到内存里,可以维护一个全局的SFTP输出流,分批写入后最后统一关闭。注意要处理线程安全问题,因为你用了@Async线程池:
// 全局维护SFTP输出流,注意线程安全 private ChannelSftp.OutputStream sftpOutputStream; private final Object streamLock = new Object(); // 初始化SFTP输出流(第一次写入前调用,或者在write方法里懒加载) private void initSftpOutputStream() throws Exception { sftpRemoteFileTemplate.execute(session -> { ChannelSftp channel = session.getClientInstance(); String remoteFilePath = contactSenderSftpProperties.getSftpSessionProperties().getBaseSftpPath() + "/contacts.csv"; try { channel.lstat(remoteFilePath); // 文件存在,以追加模式打开流 sftpOutputStream = channel.put(remoteFilePath, ChannelSftp.APPEND); } catch (SftpException e) { if (e.id == ChannelSftp.SSH_FX_NO_SUCH_FILE) { // 文件不存在,新建流 sftpOutputStream = channel.put(remoteFilePath); } else { throw new RuntimeException("SFTP流初始化失败", e); } } return null; }); } public void write(List<? extends Item> items) throws Exception { log.debug("Write {}", items); synchronized (streamLock) { // 懒加载初始化流 if (sftpOutputStream == null) { initSftpOutputStream(); } // 直接序列化到SFTP输出流 csvSerializer.serialize(sftpOutputStream, items.stream()); sftpOutputStream.flush(); // 确保数据即时写入服务器 } } // 所有数据处理完成后,必须调用这个方法关闭流 public void closeSftpStream() throws IOException { synchronized (streamLock) { if (sftpOutputStream != null) { sftpOutputStream.close(); sftpOutputStream = null; } } }
关于PipedInputStream/PipedOutputStream的疑问解答
Piped流是为单生产者、单消费者的一次性线程通信设计的,一旦流被关闭(你的try-with-resources会自动关闭PipedOutputStream),就无法再写入数据。而且你当前的实现每次write()都会新建一套Piped流,本质上是每次上传独立的小文件,自然会覆盖之前的内容。
如果想要多次追加数据,Piped流不是合适的选择——要么用上面两种SFTP直接追加的方式,要么先把数据写到本地临时文件,最后一次性上传到SFTP。
内容的提问来源于stack exchange,提问作者Tchorzyksen

