如何通过多线程优化超大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
相关产品推荐
相关产品推荐

