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

如何通过多线程优化超大CSV文件解析校验入库的处理效率

超大CSV入库多线程改造方案

这个场景最优解是生产者-消费者分层解耦模型,不要上来就搞多线程读文件——磁盘顺序读性能远高于随机读,强行拆分文件偏移量并行读不仅逻辑复杂容易出截断行的bug,实际读取效率反而比单线程顺序读低。
现有单线程逻辑慢的核心原因是三个环节(读文件、校验、入库)完全串行,读文件的时候CPU和数据库连接闲着,等数据库返回的时候磁盘IO和CPU闲着,资源完全没利用起来。

具体改造逻辑

  • 单线程做CSV读取(生产者)
    保留现有univocity CsvParser单线程迭代逻辑就行,这部分是顺序IO,单线程跑基本能打满磁盘顺序读带宽,是效率最高的读方式。不要在这个线程里加校验、入库逻辑,读到原始行直接丢进有界阻塞队列就继续读下一行,队列满了读线程自动阻塞,避免OOM。
  • CPU核心数级别的线程池做数据校验(中间处理层)
    校验是CPU密集型操作,线程数设成和CPU核心数一致就行,太多了会因为上下文切换拖慢速度。校验线程从原始行队列拿数据,跑完正则、格式校验逻辑,校验通过的行丢到待入库队列,坏数据直接归集即可。
  • 多独立连接线程做批量入库(消费者)
    这部分是单线程逻辑最大的瓶颈:JDBC单连接执行executeBatch时会阻塞等待数据库返回,这段时间其他资源全在空等。注意JDBC连接和PreparedStatement不是线程安全的,每个入库线程必须从连接池拿独立的数据库连接,自己维护自己的批量语句,攒够10002000行就执行一次批量提交。入库线程数不用太多,48个足够,太多会触发数据库锁竞争,反而降低写入速度。

关键优化点

  • 所有队列都用有界队列(比如ArrayBlockingQueue),容量设为2000~5000行就行,不要用无界队列防止内存溢出。
  • 入库必须关闭自动提交,每批执行完再手动commit,否则单条提交的性能会差一个数量级。
  • 如果用MySQL,JDBC连接串一定要加rewriteBatchedStatements=true参数,会把批量插入重写成多值INSERT语句,写入性能直接提升5~10倍,很多人改了多线程还是慢基本都是漏了这个参数。
  • 加CountDownLatch做流程同步,等读、校验、入库三个环节所有数据都处理完再释放资源,避免最后一批数据丢失。
  • 如果坏数据量很大,不要全存在内存List里,直接单独写入坏数据CSV文件,防止OOM。
  • 如果用jOOQ的loadIntoAPI,本身支持传入自定义线程池做并行批量写入,不用自己实现消费者逻辑,直接配置即可。

参考实现代码

import com.univocity.parsers.csv.CsvParser;
import com.univocity.parsers.csv.CsvParserSettings;
import javax.sql.DataSource;
import java.io.File;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.util.*;
import java.util.concurrent.*;

public class CsvParallelImporter {
    public static void main(String[] args) throws InterruptedException {
        // 基础配置
        final int COLUMN_COUNT = 2;
        final int BATCH_SIZE = 1000;
        final int QUEUE_CAPACITY = 4000;
        final int VALIDATE_THREAD = Runtime.getRuntime().availableProcessors();
        final int INSERT_THREAD = 6; // 根据数据库承载能力调整,4-8为宜
        DataSource dataSource = getYourDataSource(); // 改用连接池,不要用单连接

        // 队列定义
        BlockingQueue<String[]> rawRowQueue = new ArrayBlockingQueue<>(QUEUE_CAPACITY);
        BlockingQueue<String[]> validRowQueue = new ArrayBlockingQueue<>(QUEUE_CAPACITY);
        List<String[]> badDataList = Collections.synchronizedList(new ArrayList<>());

        // 线程池与同步计数器
        ExecutorService validatePool = Executors.newFixedThreadPool(VALIDATE_THREAD);
        ExecutorService insertPool = Executors.newFixedThreadPool(INSERT_THREAD);
        CountDownLatch parseFinish = new CountDownLatch(1);
        CountDownLatch validateFinish = new CountDownLatch(VALIDATE_THREAD);
        CountDownLatch insertFinish = new CountDownLatch(INSERT_THREAD);

        // 1. 单线程读CSV
        new Thread(() -> {
            CsvParserSettings settings = new CsvParserSettings();
            CsvParser parser = new CsvParser(settings);
            try {
                Iterator<String[]> rowIterator = parser.iterate(new File("/path/to/your.csv"), "UTF-8").iterator();
                while (rowIterator.hasNext()) {
                    rawRowQueue.put(rowIterator.next());
                }
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                parseFinish.countDown();
            }
        }).start();

        // 2. 多线程校验
        for (int i = 0; i < VALIDATE_THREAD; i++) {
            validatePool.submit(() -> {
                try {
                    while (true) {
                        if (parseFinish.getCount() == 0 && rawRowQueue.isEmpty()) break;
                        String[] row = rawRowQueue.poll(100, TimeUnit.MILLISECONDS);
                        if (row == null) continue;
                        // 替换成实际校验逻辑
                        if (row.length < 1 || !row[0].matches("^[0-9a-zA-Z]+$")) {
                            badDataList.add(row);
                            continue;
                        }
                        validRowQueue.put(row);
                    }
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    validateFinish.countDown();
                }
            });
        }

        // 3. 多线程批量入库
        for (int i = 0; i < INSERT_THREAD; i++) {
            insertPool.submit(() -> {
                try (Connection conn = dataSource.getConnection()) {
                    conn.setAutoCommit(false);
                    PreparedStatement ps = conn.prepareStatement("INSERT INTO some_table(column1, column2) VALUES (?,?)");
                    int currentBatch = 0;
                    while (true) {
                        if (validateFinish.getCount() == 0 && validRowQueue.isEmpty()) break;
                        String[] row = validRowQueue.poll(100, TimeUnit.MILLISECONDS);
                        if (row == null) continue;
                        // 字段赋值
                        for (int col = 0; col < COLUMN_COUNT; col++) {
                            ps.setObject(col + 1, col < row.length ? row[col] : null);
                        }
                        ps.addBatch();
                        currentBatch++;
                        if (currentBatch == BATCH_SIZE) {
                            ps.executeBatch();
                            conn.commit();
                            currentBatch = 0;
                        }
                    }
                    // 提交最后一批残留数据
                    if (currentBatch > 0) {
                        ps.executeBatch();
                        conn.commit();
                    }
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    insertFinish.countDown();
                }
            });
        }

        // 等待所有流程完成后关闭资源
        parseFinish.await();
        validateFinish.await();
        insertFinish.await();
        validatePool.shutdown();
        insertPool.shutdown();
    }

    private static DataSource getYourDataSource() {
        // 实现数据源初始化逻辑,记得JDBC参数加rewriteBatchedStatements=true
        return null;
    }
}

不推荐的做法

  • 不要多线程拆分读取同一个文件:除非能100%精准定位每个分块的换行边界,否则很容易把一行拆成两半,处理逻辑复杂度极高,收益还不如分层模型。
  • 不要多个入库线程共用同一个Connection:JDBC规范没有要求连接实现线程安全,共用会出现数据错乱、连接异常的问题。
  • 不要把入库线程数设得太大:数据库写入能力有上限,线程太多只会徒增锁竞争,实际测试下来4~8个线程基本能打满单表的写入上限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 18:24:45