You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

ResultSet多线程优化:解决大结果集脱敏写文件耗时过长问题

原有方案失效的核心原因

你之前的写法跑不出多线程效果,甚至处理速度极慢,根因有4个:

  • JDBC ResultSet天生不支持并发访问:Oracle JDBC驱动的ResultSet是单连接上的同步流式对象,所有读取操作都串行占用同一个数据库连接,哪怕你把逻辑丢到线程池,只要多线程访问ResultSet,就会被驱动内部的锁强制串行化,根本没有并行效果。
  • 单条任务调度开销过大:200万条记录对应200万个Runnable任务,线程池的任务创建、入队、上下文切换的开销,甚至比单条记录的脱敏处理开销还大。
  • Oracle JDBC默认配置性能极差:默认fetchSize仅为10,意味着每读10条记录就要和数据库做一次网络往返,200万条记录需要20万次网络IO,这是你10小时仅处理2万条的核心元凶之一,和多线程逻辑没有直接关系。
  • 计算和IO没有解耦:如果把脱敏计算、文件写入全混在同一个任务里,多线程写同一个文件会触发操作系统级别的文件锁竞争,最终写性能还不如单线程顺序写。
可落地的多线程优化方案

整体采用生产者-消费者三阶段解耦架构,全程做批处理,避免单条操作的 overhead:

  1. 第一阶段:单线程专属负责读ResultSet,批量拉取原始数据,所有JDBC操作只在这个线程执行,完全规避ResultSet并发问题
  2. 第二阶段:固定大小线程池批量做脱敏加密,线程数和CPU核心数对齐(加密是CPU密集型场景,线程数过多反而会因为上下文切换变慢)
  3. 第三阶段:单线程专属负责批量顺序写文件,完全规避文件锁竞争,发挥磁盘顺序写的最大性能

参考实现代码框架:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 13:12:18