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

Java多线程导出SQL:如何保证主键ID有序写入并控制内存?

多线程导出SQL文件时保证主键顺序的解决方案

我在用一款小众数据库,它不支持直接导出SQL备份,只能生成dump文件,但需求要求导出SQL,所以只能自己查数据写SQL文件。代码做了抽象处理,见谅。

import lombok.AllArgsConstructor;

import java.io.*;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class ThreadQueryWrite {

    @AllArgsConstructor
    class MyRunnable implements Runnable{

        /**
         * database tableName
         */
        String tableName;
        /**
         * The starting primary key of the query range
         */
        long startId;
        /**
         * The ending primary key of the query range
         */
        long endId;

        /**
         * File writing character stream
         */
        Writer writer;

        @Override
        public void run() {
            doQuery();
        }

        private void doQuery() {
            // resultSet = select * from tableName where id between startId and endId;
            // doWriteFile(resultSet);
        }

        private void doWriteFile(ResultSet resultSet) {
            StringBuilder sb = new StringBuilder("(");
            try {
                while(resultSet.next()){
                    // sb.append()
                }
                synchronized(writer){
                    writer.write(sb.toString());
                }
            } catch (SQLException e) {
                throw new RuntimeException(e);
            } catch (IOException e) {
                throw new RuntimeException(e);
            }
        }
    }


    public static void main(String[] args) throws IOException {
        // Assuming that the number of rows of database data
        // has been obtained is 10000000 rows
        int count = 10000000;
        String tableName = "test_table";
        String filePath = "c:\\test\\test.sql";
        ThreadQueryWrite queryWrite = new ThreadQueryWrite();
        queryWrite.doExportSql(count, tableName, filePath);
    }

    private void doExportSql(int count, String tableName, String filePath) throws IOException {
        int maxQueryCount = 3000;

        ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(Runtime.getRuntime().availableProcessors(),
                Runtime.getRuntime().availableProcessors() * 2, 5, TimeUnit.MINUTES,
                new LinkedBlockingQueue<>());

        Writer writer = new OutputStreamWriter(Files.newOutputStream(Paths.get(filePath)),
                StandardCharsets.UTF_8);
        for (int i = 1; i < count; i += maxQueryCount) {
            MyRunnable myRunnable = new MyRunnable(tableName, i, i + maxQueryCount - 1, writer);
            threadPoolExecutor.submit(myRunnable);
        }
        threadPoolExecutor.shutdown();
    }

}

问题场景

数据量大时用多线程按主键分段异步查询,希望SQL文件中主键按自增顺序连续排列,比如:

insert into table(id) values
    (1),
    ......,
    (3000);
insert into table(id) values
    (3001),
    ......,
    (6000);
......

但线程执行顺序不可控,会出现乱序(比如先写入3001-6000段,再写入1-3000段)。之前试过把所有数据加载到内存再按顺序写入,但内存压力过大,需要更优方案。


可行解决方案

方案1:分段临时文件+最终合并

每个线程查询完成后,将对应段的SQL写入单独的临时文件(按段号命名,比如temp_0001.sql、temp_0002.sql),所有线程执行完毕后,按段号顺序将临时文件合并到最终SQL文件中。

优点

  • 无锁竞争,线程独立写入,效率高
  • 内存压力小,每个线程仅缓存当前段的数据
  • 合并时严格保证顺序

代码调整示例

修改MyRunnable和导出逻辑:

@AllArgsConstructor
class MyRunnable implements Runnable{
    String tableName;
    long startId;
    long endId;
    int segmentNum;
    String tempDir;

    @Override
    public void run() {
        ResultSet resultSet = null;
        try {
            // 执行查询逻辑:select * from tableName where id between startId and endId;
            StringBuilder sb = new StringBuilder("insert into ").append(tableName).append("(id) values\n");
            while(resultSet.next()){
                sb.append("    (").append(resultSet.getLong("id")).append("),\n");
            }
            // 处理末尾逗号并添加分号
            if (sb.length() > 0 && sb.charAt(sb.length()-2) == ',') {
                sb.deleteCharAt(sb.length()-2);
            }
            sb.append(";\n");
            // 写入临时文件
            String tempFilePath = tempDir + File.separator + String.format("temp_%04d.sql", segmentNum);
            try(Writer writer = new OutputStreamWriter(Files.newOutputStream(Paths.get(tempFilePath)), StandardCharsets.UTF_8)){
                writer.write(sb.toString());
            }
        } catch (SQLException | IOException e) {
            throw new RuntimeException(e);
        } finally {
            // 关闭ResultSet等资源
            if (resultSet != null) {
                try {
                    resultSet.close();
                } catch (SQLException ignored) {}
            }
        }
    }
}

private void doExportSql(int count, String tableName, String filePath) throws IOException, InterruptedException {
    int maxQueryCount = 3000;
    String tempDir = "c:\\test\\temp_sql_segments";
    // 创建临时目录
    Files.createDirectories(Paths.get(tempDir));

    ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(
            Runtime.getRuntime().availableProcessors(),
            Runtime.getRuntime().availableProcessors() * 2,
            5, TimeUnit.MINUTES,
            new LinkedBlockingQueue<>()
    );

    int totalSegments = (count + maxQueryCount - 1) / maxQueryCount;
    for (int i = 1; i <= count; i += maxQueryCount) {
        int segmentNum = (i - 1) / maxQueryCount + 1;
        long endId = Math.min(i + maxQueryCount - 1, count);
        MyRunnable myRunnable = new MyRunnable(tableName, i, endId, segmentNum, tempDir);
        threadPoolExecutor.submit(myRunnable);
    }

    // 等待所有线程完成
    threadPoolExecutor.shutdown();
    if (!threadPoolExecutor.awaitTermination(1, TimeUnit.HOURS)) {
        threadPoolExecutor.shutdownNow();
        throw new RuntimeException("导出任务超时");
    }

    // 合并临时文件到最终SQL
    try(Writer finalWriter = new OutputStreamWriter(Files.newOutputStream(Paths.get(filePath)), StandardCharsets.UTF_8)){
        for (int i = 1; i <= totalSegments; i++) {
            String tempFilePath = tempDir + File.separator + String.format("temp_%04d.sql", i);
            String segmentContent = Files.readString(Paths.get(tempFilePath), StandardCharsets.UTF_8);
            finalWriter.write(segmentContent);
            // 删除临时文件
            Files.delete(Paths.get(tempFilePath));
        }
    }
    // 删除临时目录
    Files.delete(Paths.get(tempDir));
}

方案2:有序阻塞队列+单独写文件线程

  1. 定义包含段号和SQL内容的结果对象,每个线程查询完成后将对象放入阻塞队列。
  2. 启动单独的写文件线程,维护当前需要写入的段号(从1开始):
    • 从队列中取出结果,若段号等于当前待写段号则直接写入,否则暂存到Map中
    • 写入后段号递增,检查Map中是否存在下一段的数据,存在则继续写入

优点

  • 无需临时文件,逻辑更紧凑
  • 内存压力可控,仅缓存少量待写入的段数据

内容的提问来源于stack exchange,提问作者HuaJFrame

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 12:04:51