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

Java多线程实现文件批量处理与H2数据库快速插入的最优方案

实现高并发H2数据插入的完整方案(适合多线程新手)

Hey there! Let's break this down step by step since you're new to multi-threaded coding—we'll hit that 2-second insertion goal with H2, no sweat. First, let's lock in the core principles that make this fast:

  • Batch inserts are non-negotiable (single-row inserts will never hit your timing target)
  • A dynamically sized thread pool that scales with your record count to avoid wasted resources
  • Minimizing file IO competition by pre-splitting data into batches for threads

1. 必备依赖(Maven示例)

We'll use HikariCP (the fastest JDBC connection pool) and the official H2 driver:

<dependencies>
    <!-- H2 Database Driver -->
    <dependency>
        <groupId>com.h2database</groupId>
        <artifactId>h2</artifactId>
        <version>2.2.224</version>
    </dependency>
    <!-- HikariCP Connection Pool -->
    <dependency>
        <groupId>com.zaxxer</groupId>
        <artifactId>HikariCP</artifactId>
        <version>5.0.1</version>
    </dependency>
</dependencies>

2. 核心实现步骤

2.1 先创建H2表结构

Assuming your 4 columns are col1 to col4, run this SQL first (you can execute it once in your app or via H2 console):

CREATE TABLE IF NOT EXISTS data_table (
    col1 VARCHAR(255),
    col2 VARCHAR(255),
    col3 VARCHAR(255),
    col4 VARCHAR(255)
);

2.2 配置高性能H2连接池

HikariCP is optimized for speed, and we'll tweak H2 settings to maximize batch insert performance:

import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;

public class H2DataSourceUtil {
    private static HikariDataSource dataSource;

    static {
        HikariConfig config = new HikariConfig();
        // Use in-memory mode for maximum speed (swap to file mode if you need persistence)
        config.setJdbcUrl("jdbc:h2:mem:test;DB_CLOSE_DELAY=-1;MODE=MySQL;CACHE_SIZE=65536");
        config.setUsername("sa");
        config.setPassword("");
        config.setMaximumPoolSize(20); // Base max size, we'll adjust dynamically
        config.setMinimumIdle(5);
        config.setConnectionTimeout(3000);
        dataSource = new HikariDataSource(config);
    }

    public static HikariDataSource getDataSource() {
        return dataSource;
    }

    // Match connection pool size to thread pool size to avoid connection shortages
    public static void adjustPoolSize(int newMaxSize) {
        dataSource.setMaximumPoolSize(Math.min(newMaxSize, 50)); // Cap at 50 to avoid overloading H2
    }
}

2.3 动态线程池实现

We'll scale the thread pool based on your total record count to balance speed and resource usage:

import java.util.concurrent.*;

public class DynamicThreadPool {
    private static ThreadPoolExecutor executor;

    public static ThreadPoolExecutor getExecutor(int recordCount) {
        int cpuCores = Runtime.getRuntime().availableProcessors();
        int corePoolSize;
        int maxPoolSize;

        // Scale based on record volume
        if (recordCount < 100000) {
            corePoolSize = cpuCores;
            maxPoolSize = cpuCores * 2;
        } else if (recordCount < 200000) {
            corePoolSize = cpuCores * 2;
            maxPoolSize = cpuCores * 4;
        } else {
            corePoolSize = cpuCores * 4;
            maxPoolSize = cpuCores * 8;
        }

        // Sync connection pool size with thread pool
        H2DataSourceUtil.adjustPoolSize(maxPoolSize);

        executor = new ThreadPoolExecutor(
                corePoolSize,
                maxPoolSize,
                60L,
                TimeUnit.SECONDS,
                new LinkedBlockingQueue<>(1000), // Spawn new threads when queue is full
                Executors.defaultThreadFactory(),
                new ThreadPoolExecutor.CallerRunsPolicy() // Fall back to main thread if queue overflows to avoid data loss
        );
        return executor;
    }

    public static void shutdown() {
        if (executor != null) {
            executor.shutdown();
            try {
                if (!executor.awaitTermination(10, TimeUnit.SECONDS)) {
                    executor.shutdownNow();
                }
            } catch (InterruptedException e) {
                executor.shutdownNow();
            }
        }
    }
}

2.4 批量插入任务(每个线程处理一批数据)

Each thread gets its own JDBC connection (never share connections—they're not thread-safe!) and handles a batch of records:

import java.sql.Connection;
import java.sql.PreparedStatement;
import java.util.List;

public class InsertTask implements Runnable {
    private static final String INSERT_SQL = "INSERT INTO data_table (col1, col2, col3, col4) VALUES (?, ?, ?, ?)";
    private List<String> lines;

    public InsertTask(List<String> lines) {
        this.lines = lines;
    }

    @Override
    public void run() {
        // Try-with-resources auto-closes connection/statement
        try (Connection conn = H2DataSourceUtil.getDataSource().getConnection();
             PreparedStatement pstmt = conn.prepareStatement(INSERT_SQL)) {

            conn.setAutoCommit(false); // Disable auto-commit for batch efficiency
            for (String line : lines) {
                // Split into max 4 columns (handles spaces inside columns)
                String[] cols = line.split(" ", 4);
                if (cols.length == 4) {
                    pstmt.setString(1, cols[0]);
                    pstmt.setString(2, cols[1]);
                    pstmt.setString(3, cols[2]);
                    pstmt.setString(4, cols[3]);
                    pstmt.addBatch();
                }
            }
            pstmt.executeBatch();
            conn.commit();

        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

2.5 主程序:读取文件、分批次、执行插入

import java.io.BufferedReader;
import java.io.FileReader;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ThreadPoolExecutor;

public class MainInsertApp {
    private static final int BATCH_SIZE = 1000; // Adjust this based on your testing (1000-2000 works well)

    public static void main(String[] args) throws Exception {
        long startTime = System.currentTimeMillis();

        // 1. Read all lines into memory (50k-300k lines is trivial for modern RAM)
        List<String> allLines = new ArrayList<>();
        try (BufferedReader br = new BufferedReader(new FileReader("your-data-file.txt"))) {
            String line;
            while ((line = br.readLine()) != null) {
                if (!line.trim().isEmpty()) {
                    allLines.add(line);
                }
            }
        }

        int totalRecords = allLines.size();
        System.out.println("Total records to insert: " + totalRecords);

        // 2. Get dynamically sized thread pool
        ThreadPoolExecutor executor = DynamicThreadPool.getExecutor(totalRecords);

        // 3. Split data into batches and submit tasks
        for (int i = 0; i < totalRecords; i += BATCH_SIZE) {
            int endIndex = Math.min(i + BATCH_SIZE, totalRecords);
            List<String> batch = allLines.subList(i, endIndex);
            executor.submit(new InsertTask(batch));
        }

        // 4. Wait for all tasks to finish
        DynamicThreadPool.shutdown();

        long endTime = System.currentTimeMillis();
        System.out.println("Insert completed in " + (endTime - startTime) + " ms");
    }
}

3. Critical Optimizations to Hit 2 Seconds

  • H2 In-Memory Mode: jdbc:h2:mem:test is drastically faster than file-based storage (swap to file mode only if you need persistence)
  • Batch Commits: Disabling auto-commit and using executeBatch() cuts down on database round-trips
  • Dynamic Scaling: Matching thread/connection pool size to record volume avoids underutilization or context-switching overhead
  • Pre-Read File: Loading all lines into memory first prevents multiple threads from competing for file IO

4. Newbie-Friendly Notes

  • Never share JDBC connections between threads—each task should grab its own from the pool
  • The split(" ", 4) trick ensures you don't break columns that contain spaces
  • Test with smaller datasets first to tweak BATCH_SIZE (1000 is a safe starting point)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:32:30