如何通过Spring Integration实现SFTP流式分块读取并对接Spring Batch?
解决方案
Spring Integration对SFTP分块读取的支持
Spring Integration的SFTP组件完全支持分块读取远程文件,无需一次性把整个文件加载到堆内存:
- 放弃默认
SftpStreamingMessageSource的一次性读取逻辑,改用SftpRemoteFileTemplate的get()方法获取远程文件的输入流。这个流基于JSch实现,支持skip()和分块读取操作,不会将文件下载到本地临时目录。 - 你可以完全控制读取的块大小和起始位置,从根源避免大文件占用过多内存。
对接Spring Batch的无本地存储方案
核心思路是自定义Spring Batch的ItemReader,结合Spring Integration的SFTP组件实现流式分块读取,同时通过ItemStream接口实现失败恢复:
1. 配置SFTP基础组件
先配置SFTP会话工厂和远程文件模板,用于连接SFTP服务器并操作文件:
@Bean public SftpSessionFactory sftpSessionFactory() { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); factory.setHost("你的SFTP主机"); factory.setPort(22); factory.setUser("用户名"); factory.setPassword("密码"); // 若使用密钥认证,替换为setPrivateKey等对应配置 return factory; } @Bean public SftpRemoteFileTemplate sftpRemoteFileTemplate(SftpSessionFactory sessionFactory) { return new SftpRemoteFileTemplate(sessionFactory); }
2. 自定义支持恢复的分块ItemReader
实现ItemReader和ItemStream接口,在open()方法中从Spring Batch元数据恢复读取偏移量,read()方法分块读取数据,update()方法保存当前读取位置:
public class SftpChunkItemReader implements ItemReader<String>, ItemStream { private final SftpRemoteFileTemplate sftpTemplate; private final String remoteFilePath; private final int chunkSize; // 自定义文件读取块大小 private BufferedReader reader; private long currentOffset = 0; private static final String OFFSET_CONTEXT_KEY = "sftp_file_read_offset"; public SftpChunkItemReader(SftpRemoteFileTemplate sftpTemplate, String remoteFilePath, int chunkSize) { this.sftpTemplate = sftpTemplate; this.remoteFilePath = remoteFilePath; this.chunkSize = chunkSize; } @Override public void open(ExecutionContext executionContext) { // 从Spring Batch元数据恢复上次读取的位置 if (executionContext.containsKey(OFFSET_CONTEXT_KEY)) { currentOffset = executionContext.getLong(OFFSET_CONTEXT_KEY); } // 获取远程文件流并跳转到指定偏移量 InputStream remoteStream = sftpTemplate.get(remoteFilePath, is -> { try { is.skip(currentOffset); } catch (IOException e) { throw new RuntimeException("Failed to skip to offset: " + currentOffset, e); } return is; }); this.reader = new BufferedReader(new InputStreamReader(remoteStream)); } @Override public String read() throws IOException { StringBuilder chunk = new StringBuilder(); char[] buffer = new char[chunkSize]; int readCount = reader.read(buffer); if (readCount == -1) { close(); return null; } chunk.append(buffer, 0, readCount); currentOffset += readCount; return chunk.toString(); } @Override public void update(ExecutionContext executionContext) { // 保存当前偏移量到Spring Batch元数据 executionContext.putLong(OFFSET_CONTEXT_KEY, currentOffset); } @Override public void close() { try { if (reader != null) { reader.close(); } } catch (IOException e) { throw new RuntimeException("Failed to close SFTP stream", e); } } }
3. 配置Spring Batch Job
将自定义的SftpChunkItemReader注入到Step中,结合你的业务ItemProcessor和ItemWriter即可:
@Bean public Step sftpProcessingStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, SftpChunkItemReader sftpReader, ItemProcessor<String, YourDto> processor, ItemWriter<YourDto> writer) { return new StepBuilder("sftpProcessingStep", jobRepository) .<String, YourDto>chunk(1, transactionManager) // 此处为Batch处理块,与文件读取块区分开 .reader(sftpReader) .processor(processor) .writer(writer) .build(); } @Bean public SftpChunkItemReader sftpChunkItemReader(SftpRemoteFileTemplate sftpTemplate) { return new SftpChunkItemReader(sftpTemplate, "/remote/path/your-file.txt", 1024); // 1KB文件读取块 }
关键注意事项
- 确保SFTP服务器支持文件随机访问(主流SFTP服务器均支持,JSch的
ChannelSftp默认支持skip()操作)。 - 读取块大小按需调整:文本文件可改为按行读取(将
read()方法替换为reader.readLine()或读取多行),二进制文件按固定字节块处理。 - 应用崩溃后,只需移除Spring Batch元数据存储中对应JobExecution的条目,重启后Job会从上次失败的偏移量处继续读取。
内容的提问来源于stack exchange,提问作者Harpreet Saund
相关产品推荐
相关产品推荐

