Java使用JSch实现数据库批量数据流式压缩上传至SFTP服务器
问题翻译
我需要用Java从数据库加载批量数据,并将其保存到远程SFTP服务器的单个bzip2压缩文件中,目前使用JSch作为SFTP客户端。请问:
- 是否可以动态压缩每个数据批次(chunk)并拼接成最终文件?
- 能否使用BufferedInputStream,在每次迭代时压缩缓冲区并传递给JSch的
put(InputStream src, String dst)方法? - 由于数据量较大,我不想将所有数据存储在内存或本地文件中,bzip2格式是否支持这种流式操作?
解决方案
核心结论
bzip2完全支持流式压缩,不需要把所有数据加载到内存或写入本地文件就能完成远程SFTP上传。你可以通过自定义流式处理逻辑,结合JSch的put方法实现需求。
具体实现思路
流式压缩+SFTP上传的核心逻辑
不要直接用BufferedInputStream逐个批次压缩后拼接,而是用**管道流(PipedInputStream/PipedOutputStream)**配合多线程实现:- 主线程:将
PipedInputStream传给JSch的put方法,让SFTP客户端从这个流读取数据并上传。 - 子线程:从数据库批量读取数据,将原始数据写入
BZip2CompressorOutputStream(该流连接到PipedOutputStream),实时生成bzip2压缩字节。
- 主线程:将
关键依赖
推荐使用org.apache.commons:commons-compress库来处理bzip2流式压缩,它提供了可靠的BZip2CompressorOutputStream实现。
代码示例
import com.jcraft.jsch.ChannelSftp; import com.jcraft.jsch.JSch; import com.jcraft.jsch.Session; import org.apache.commons.compress.compressors.bzip2.BZip2CompressorOutputStream; import java.io.PipedInputStream; import java.io.PipedOutputStream; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; public class SftpBzip2StreamUpload { public static void main(String[] args) throws Exception { // 1. 初始化SFTP会话和通道 JSch jsch = new JSch(); Session session = jsch.getSession("sftp-user", "sftp-host", 22); session.setPassword("sftp-password"); session.setConfig("StrictHostKeyChecking", "no"); session.connect(); ChannelSftp channel = (ChannelSftp) session.openChannel("sftp"); channel.connect(); // 2. 创建管道流 PipedInputStream in = new PipedInputStream(); PipedOutputStream out = new PipedOutputStream(in); // 3. 启动子线程负责读取数据库数据并压缩 Thread compressThread = new Thread(() -> { try (Connection conn = getDbConnection(); // 自行实现数据库连接获取 PreparedStatement stmt = conn.prepareStatement("SELECT data_column FROM large_table"); ResultSet rs = stmt.executeQuery(); BZip2CompressorOutputStream bzipOut = new BZip2CompressorOutputStream(out)) { byte[] buffer = new byte[8192]; while (rs.next()) { // 读取单条数据(或批量读取,根据你的需求调整) String data = rs.getString("data_column"); byte[] rawData = data.getBytes("UTF-8"); // 将原始数据写入压缩流,自动生成压缩字节到管道 bzipOut.write(rawData); // 可选:每批次后刷新,避免数据在缓冲区堆积 bzipOut.flush(); } // 完成压缩,关闭输出流 bzipOut.finish(); } catch (Exception e) { e.printStackTrace(); try { out.close(); } catch (Exception ex) { ex.printStackTrace(); } } }); compressThread.start(); // 4. 主线程用JSch上传管道流中的压缩数据 try { channel.put(in, "/remote/path/output.bz2"); } finally { in.close(); channel.disconnect(); session.disconnect(); } // 等待压缩线程完成 compressThread.join(); } // 自行实现数据库连接方法 private static Connection getDbConnection() throws Exception { // 示例:返回你的数据库连接 return null; } }
注意事项
- 线程安全:管道流的输入输出端必须分别在不同线程中使用,否则会导致死锁。
- 批次读取优化:可以调整数据库查询的批量读取逻辑(比如用
rs.setFetchSize()),减少数据库交互次数,提升效率。 - 异常处理:要确保在任何异常情况下,流和SFTP会话都能正确关闭,避免资源泄漏。
- 压缩参数调整:
BZip2CompressorOutputStream支持设置压缩级别,可根据速度和压缩比需求调整。
内容的提问来源于stack exchange,提问作者200OK
相关产品推荐
相关产品推荐

