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:有序阻塞队列+单独写文件线程
- 定义包含段号和SQL内容的结果对象,每个线程查询完成后将对象放入阻塞队列。
- 启动单独的写文件线程,维护当前需要写入的段号(从1开始):
- 从队列中取出结果,若段号等于当前待写段号则直接写入,否则暂存到Map中
- 写入后段号递增,检查Map中是否存在下一段的数据,存在则继续写入
优点
- 无需临时文件,逻辑更紧凑
- 内存压力可控,仅缓存少量待写入的段数据
内容的提问来源于stack exchange,提问作者HuaJFrame
相关产品推荐
相关产品推荐

