ResultSet多线程优化:解决大结果集脱敏写文件耗时过长问题
原有方案失效的核心原因
你之前的写法跑不出多线程效果,甚至处理速度极慢,根因有4个:
- JDBC ResultSet天生不支持并发访问:Oracle JDBC驱动的ResultSet是单连接上的同步流式对象,所有读取操作都串行占用同一个数据库连接,哪怕你把逻辑丢到线程池,只要多线程访问ResultSet,就会被驱动内部的锁强制串行化,根本没有并行效果。
- 单条任务调度开销过大:200万条记录对应200万个Runnable任务,线程池的任务创建、入队、上下文切换的开销,甚至比单条记录的脱敏处理开销还大。
- Oracle JDBC默认配置性能极差:默认
fetchSize仅为10,意味着每读10条记录就要和数据库做一次网络往返,200万条记录需要20万次网络IO,这是你10小时仅处理2万条的核心元凶之一,和多线程逻辑没有直接关系。 - 计算和IO没有解耦:如果把脱敏计算、文件写入全混在同一个任务里,多线程写同一个文件会触发操作系统级别的文件锁竞争,最终写性能还不如单线程顺序写。
可落地的多线程优化方案
整体采用生产者-消费者三阶段解耦架构,全程做批处理,避免单条操作的 overhead:
- 第一阶段:单线程专属负责读ResultSet,批量拉取原始数据,所有JDBC操作只在这个线程执行,完全规避ResultSet并发问题
- 第二阶段:固定大小线程池批量做脱敏加密,线程数和CPU核心数对齐(加密是CPU密集型场景,线程数过多反而会因为上下文切换变慢)
- 第三阶段:单线程专属负责批量顺序写文件,完全规避文件锁竞争,发挥磁盘顺序写的最大性能
参考实现代码框架:
import java.io.BufferedWriter; import java.io.FileWriter; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; public class OracleCursorExporter { // 配置参数,根据实际服务器内存、CPU调整 private static final int BATCH_SIZE = 1000; private static final int WORKER_THREAD_NUM = Runtime.getRuntime().availableProcessors(); private static final int QUEUE_CAPACITY = 20; private static final int COLUMN_COUNT = 400; private static final List<RawRecord> READ_FINISH_FLAG = Collections.emptyList(); private static final List<MaskedRecord> PROCESS_FINISH_FLAG = Collections.emptyList(); // 预加载好400列的脱敏配置,不要每条记录临时判断 private static final boolean[] NEED_MASK_COL = MaskConfig.loadAllColumnMaskFlag(); // 用ThreadLocal给每个工作线程缓存初始化好的加密Cipher,避免重复创建开销 private static final ThreadLocal<Cipher> CIPHER_CACHE = ThreadLocal.withInitial(() -> EncryptUtil.initCipher()); public void doExport(Statement stmt, String outputPath) throws Exception { // 提前设置JDBC读取参数,这步对性能影响极大 stmt.setFetchSize(5000); stmt.setFetchDirection(ResultSet.FETCH_FORWARD); BlockingQueue<List<RawRecord>> processQueue = new ArrayBlockingQueue<>(QUEUE_CAPACITY); BlockingQueue<List<MaskedRecord>> writeQueue = new ArrayBlockingQueue<>(QUEUE_CAPACITY); // 1. 生产者线程:唯一负责读取ResultSet,批量拉取原始数据 Thread readThread = new Thread(() -> { List<RawRecord> batch = new ArrayList<>(BATCH_SIZE); try (ResultSet rs = stmt.executeQuery()) { while (rs.next()) { RawRecord record = new RawRecord(COLUMN_COUNT); // 只做最简单的字段取值,不做任何脱敏逻辑 for (int i = 1; i <= COLUMN_COUNT; i++) { record.setField(i-1, rs.getObject(i)); } batch.add(record); if (batch.size() >= BATCH_SIZE) { processQueue.put(batch); batch = new ArrayList<>(BATCH_SIZE); } } if (!batch.isEmpty()) { processQueue.put(batch); } // 放入结束标记 processQueue.put(READ_FINISH_FLAG); } catch (Exception e) { throw new RuntimeException("读取游标失败", e); } }); // 2. 处理线程池:批量做脱敏加密 ExecutorService processPool = Executors.newFixedThreadPool(WORKER_THREAD_NUM); for (int i = 0; i < WORKER_THREAD_NUM; i++) { processPool.execute(() -> { try { while (true) { List<RawRecord> rawBatch = processQueue.take(); if (rawBatch == READ_FINISH_FLAG) { writeQueue.put(PROCESS_FINISH_FLAG); break; } List<MaskedRecord> maskedBatch = new ArrayList<>(rawBatch.size()); Cipher cipher = CIPHER_CACHE.get(); for (RawRecord raw : rawBatch) { MaskedRecord masked = new MaskedRecord(COLUMN_COUNT); for (int col = 0; col < COLUMN_COUNT; col++) { Object val = raw.getField(col); if (NEED_MASK_COL[col] && val != null) { masked.setField(col, EncryptUtil.encrypt(cipher, val)); } else { masked.setField(col, val); } } maskedBatch.add(masked); } writeQueue.put(maskedBatch); } } catch (Exception e) { throw new RuntimeException("脱敏处理失败", e); } }); } // 3. 写入线程:单线程批量顺序写文件,性能最高 Thread writeThread = new Thread(() -> { // 设置1M写缓冲区,不要频繁刷盘 try (BufferedWriter writer = new BufferedWriter(new FileWriter(outputPath), 1024*1024)) { int finishCnt = 0; while (true) { List<MaskedRecord> batch = writeQueue.take(); if (batch == PROCESS_FINISH_FLAG) { finishCnt++; if (finishCnt >= WORKER_THREAD_NUM) break; continue; } for (MaskedRecord record : batch) { writer.write(record.toLine()); writer.newLine(); } } } catch (Exception e) { throw new RuntimeException("写文件失败", e); } }); // 启动所有任务等待完成 readThread.start(); writeThread.start(); readThread.join(); processPool.shutdown(); processPool.awaitTermination(1, TimeUnit.DAYS); writeThread.join(); } }
关键避坑点
- 绝对不要在工作线程里访问ResultSet,所有JDBC相关的
rs.next()、rs.getObject()操作必须全部放在读线程里执行,否则会出现数据错乱、驱动抛异常、锁串行化等问题。 - 加密用的Cipher对象不要每条记录新建,也不要多线程共用,用ThreadLocal缓存每个线程专属的初始化好的Cipher对象,能减少90%以上的加密初始化开销。
- 不要多线程写同一个文本文件,多线程写文件要么出现行内容错乱,要么因为锁竞争导致性能远低于单线程顺序写。
- 用有界阻塞队列控制内存占用,不要一次性把200万条记录全读到内存里,队列容量控制在20左右、每批1000条的情况下,内存中仅驻留2万条左右记录,不会出现OOM,同时能保证各阶段线程不会空等。
- 不要每写一条记录就刷盘,利用BufferedWriter的缓冲区做批量写,最后关流时统一刷盘即可,频繁刷盘的随机IO开销会比顺序写高2~3个数量级。
内容的提问来源于stack exchange,提问作者slowjava
相关产品推荐
相关产品推荐

