刚看到一个很有意思的搜索词:"cvs导入多连接池"。第一眼我以为是版本管理工具CVS,再一看上下文全是"导入csv文件""excel导入数据库"这类需求,瞬间明白了——这大概率是CSV的笔误。不过这个误打误撞的组合词,反而精准戳中了一个在企业级开发里非常现实的痛点:当CSV文件体积变大、数据量从几万行涨到几百万行时,单线程逐行导入数据库的做法直接就废了,慢到怀疑人生。
今天我想围绕这个主题,把"多线程 + 多连接池并行导入CSV"这套方案从头到尾聊透。咱们不整虚的,直接讲清楚原理、代码怎么写、参数怎么调、坑在哪,让你看完能直接拿去用。
1. 先搞懂"多连接池"到底在解决什么问题
1.1 单连接导入为什么慢到让你抓狂
先回忆一下最原始的导入方式:打开一个数据库连接,用INSERT INTO一条一条地往表里插数据。这个方案的性能瓶颈非常明显,总结下来就三个字:太串行。
一次插入操作,在数据库层面至少经历这么几步:客户端发送SQL语句、服务端解析SQL、执行计划生成、锁竞争、事务日志写入、最终落盘。如果是在远程数据库上,还得叠加每一轮的网络往返延迟。假设一次插入需要5毫秒,1万条数据就是50秒,100万条数据就是5000秒,一个多小时过去了,业务早就炸了。
更麻烦的是事务开销。如果每条插入都单独提交一次事务,那每一笔都要额外承担一次fsync刷盘和事务协调的代价,性能雪上加霜。
所以单连接慢的根本原因是:所有操作串行排队,资源利用率极低,数据库的吞吐能力被完全浪费掉了。
1.2 连接池、线程池、多连接池到底是什么关系
很多人一谈到"多连接池"就懵,其实拆开看不复杂。
- 连接池:维护一组现成的数据库连接的容器。它解决的是"反复创建/销毁连接太慢"的问题,连接用完归还而不是关闭。
- 线程池:维护一组工作线程的容器。它解决的是"频繁创建/销毁线程开销大"以及"无限制并发导致资源耗尽"的问题。
- 多连接池:在导入场景下更常见的理解,是给导入任务单独建一个专属连接池(或者多个连接池),不让导入任务和线上业务互相挤占连接资源。另一种理解是多个线程各从连接池中取连接,实现真正的并行写入。
打个生活化的比方:单连接导入就像只有一个收银台的超市,所有顾客排一队,后面的人只能干等。多连接池就是开了多个收银台,多支队伍同时结账,整体吞吐量自然就上去了。
1.3 什么场景才值得上"多连接池"
不是所有CSV导入都需要这么复杂的方案。我自己的经验是这么判断的:
- 数据量在1万行以内:单连接批量提交就够用了,别折腾。
- 数据量在1万到50万行:单连接 + 批量提交(
addBatch+executeBatch)基本能扛住,但已经有明显延迟。 - 数据量在50万行以上:就必须上多线程 + 多连接池了,否则等不起。
- 文件超大(超过500MB):除了并行写库,还得考虑流式读取和分片,否则内存先爆。
另外还有一个信号:如果你发现导入过程中,线上业务查询明显变慢,说明你的导入连接把数据库资源吃满了。这时候把导入任务隔离到独立连接池,是个非常明智的做法。
2. 整体设计思路:一条CSV是怎么被"拆"进数据库的
2.1 六阶段流水线架构
我之前在一个数据迁移项目里踩过很多坑,后来沉淀出一套比较通用的并行导入流水线,分成六个阶段:
- 文件扫描与规划:读取文件基本信息,确定行数、预估分片数量。
- 流式读取解析:不一次性把CSV全load进内存,而是逐行读取、逐行分发。
- 数据校验与清洗:处理空值、类型转换、长度校验、编码修正。
- 分片打包:把校验通过的行按固定大小(如每批1000行)打包成任务块。
- 并行写库:线程池领取任务块,从独立连接池获取连接,批量写入。
- 结果汇总与重试:统计成功/失败行数,失败块进入重试队列。
这六个阶段里,最容易做错的是第4和第5步。很多人一上来就搞"每行一个线程",结果线程数爆炸、数据库连接池被瞬间打满、锁竞争剧烈,性能反而比单线程还差。
正确的做法是:按"批"而不是按"行"做并行单元。每批500到2000行,既能让数据库批量执行语句,又不至于事务太长导致锁范围过大。
2.2 文件分片策略:怎么拆才合理
分片方式直接决定导入效率,我试过三种方案,各有优劣:
- 按行数均匀分片:先统计总行数,除以期望的分片数,得到每片行数,然后各线程读取指定行区间。逻辑简单,但需要先遍历一遍文件统计行数,而且行长度不均匀时会导致负载不均衡。
- 按字节偏移分片:把文件按大小均分成N段,从各段起始位置找最近的行尾并自行处理首尾行。速度快,适合超大文件,但要自己处理跨段的半行,容易出bug。
- 按读取队列动态分发:一个主线程负责流式读取CSV,读到的行推入有界队列,多个工作线程从队列取数据进行分批。实现稍微复杂,但内存可控、负载均衡最好,因为不会出现"某些线程累死,某些线程闲死"的情况。
我最终推荐第三种方案。原因很简单:前两种方案虽然实现简单,但在文件行长度差异大(比如备注字段内容长度相差十几倍)时,分片之间严重倾斜,快的线程干完没事干,慢的线程累成狗。动态分发天然解决了负载不均衡的问题,队列的背压机制还能防止内存被撑爆。
2.3 线程池、连接池参数映射关系
这里有一个非常常见的误区:线程池大小 = 连接池大小。
我遇到过一个案例,开发人员把连接池最大连接数设为8,线程池核心线程数却设成了32。结果32个线程抢8个连接,一半线程在阻塞等待连接,CPU和数据库都没干到满负荷,整体吞吐反而上不去。
一个稳妥的经验公式:
对于单文件导入任务:连接池最大连接数 = 线程池并发线程数。每个工作线程执行写库操作时,都必须先拿到一个连接,所以连接数少于线程数只会徒增等待。
线程数的计算公式则要结合机器核数和数据库能力:
推荐起始值:CPU核数 × 2,然后根据实测上下调整。如果数据库是远程主从架构,带宽、连接数限制也要计入。
比如4核机器,起始开8个线程,连接池上限也设为8,先跑一轮看耗时,再逐步加到12、16,观察数据库CPU和锁等待指标的拐点。
3. 核心实现:一个可直接参考的并行导入骨架
3.1 工程结构和依赖准备
我用Java来写示例,这套思路同样适用于Python、Go、C#,核心逻辑是一样的。工程里需要的东西:
- JDK 8+(建议17,虚拟线程方案更优雅)
- 数据库驱动:MySQL Connector/J(或者对应的PostgreSQL、Oracle驱动)
- 连接池:HikariCP(目前综合表现最强的连接池,没有之一)
- CSV解析:直接
BufferedReader按行读就行,不引入复杂依赖
Maven依赖非常简单:
<dependency> <groupId>com.zaxxer</groupId> <artifactId>HikariCP</artifactId> <version>5.0.1</version> </dependency> <dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>8.0.33</version> </dependency>3.2 独立连接池的创建
先强调一件事:导入用的连接池务必和业务连接池分开。原因有两个:
一是隔离。导入任务经常是低频但高强度的操作,如果和线上业务共用连接池,导入瞬间会把连接全部抢走,线上查询直接超时报警。
二是参数差异。导入场景希望连接池大、事务自动提交关闭、批量写入优化参数打开,而业务连接池通常不需要这些设置。混在一起配置会很别扭。
public HikariDataSource buildImportDataSource(String jdbcUrl, String username, String password, int maxPoolSize) { HikariConfig config = new HikariConfig(); config.setJdbcUrl(jdbcUrl); config.setUsername(username); config.setPassword(password); config.setMaximumPoolSize(maxPoolSize); config.setMinimumIdle(1); config.setConnectionTimeout(30_000); config.setAutoCommit(false); // 批量导入场景,交给代码控制事务边界 config.setPoolName("import-pool"); // 这两个参数是针对MySQL批量写入的核心优化 config.addDataSourceProperty("rewriteBatchedStatements", "true"); config.addDataSourceProperty("useServerPrepStmts", "true"); return new HikariDataSource(config); }这里最关键的就是rewriteBatchedStatements=true。没有这个参数,executeBatch在MySQL驱动层还是逐条发送SQL,性能提升极其有限。开了之后,驱动会把多条插入语句重写成一条多值插入语句,性能有质的飞跃。
3.3 流式读取 + 动态分发实现
动态分发最经典的做法是用BlockingQueue作为读线程和工作线程之间的缓冲。
public void parallelImport(String filePath, DataSource dataSource, int threadCount, int batchSize) throws Exception { // 有界队列,防止读太快把内存打爆 BlockingQueue<List<String[]>> queue = new ArrayBlockingQueue<>(threadCount * 2); AtomicLong successCount = new AtomicLong(0); AtomicLong failCount = new AtomicLong(0); // 工作线程池 ExecutorService workers = Executors.newFixedThreadPool(threadCount); CountDownLatch latch = new CountDownLatch(threadCount); // 启动N个工作线程 for (int i = 0; i < threadCount; i++) { workers.execute(() -> { try { while (true) { List<String[]> batch = queue.poll(5, TimeUnit.SECONDS); if (batch == null) { // 没有新任务且读线程已结束,退出 if (queue.isEmpty() && readFinished.get()) { break; } continue; } int inserted = writeBatch(dataSource, batch); successCount.addAndGet(inserted); } } catch (Exception e) { failCount.incrementAndGet(); log.error("工作线程异常", e); } finally { latch.countDown(); } }); } // 主线程流式读取并分发 try (BufferedReader reader = new BufferedReader(new InputStreamReader(new FileInputStream(filePath), StandardCharsets.UTF_8))) { String line = reader.readLine(); // 跳过表头 List<String[]> batch = new ArrayList<>(batchSize); while ((line = reader.readLine()) != null) { String[] fields = parseCsvLine(line); if (fields == null) continue; batch.add(fields); if (batch.size() >= batchSize) { queue.put(batch); // 队列满时自动阻塞,形成背压 batch = new ArrayList<>(batchSize); } } if (!batch.isEmpty()) { queue.put(batch); } readFinished.set(true); } // 等待所有线程结束 latch.await(5, TimeUnit.MINUTES); workers.shutdown(); }队列的put方法是阻塞的,意思就是读线程发现队列满了就自动停下来等。这个机制特别重要:它天然地把"读文件速度"和"写数据库速度"做了适配。如果写库慢,读线程就自动变慢,内存占用始终有上限。
3.4 批量写入核心代码
写入方法是性能调优的重头戏。
private int writeBatch(DataSource dataSource, List<String[]> batch) { String sql = "INSERT INTO target_table (col1, col2, col3) VALUES (?, ?, ?)"; try (Connection conn = dataSource.getConnection(); PreparedStatement ps = conn.prepareStatement(sql)) { for (String[] row : batch) { ps.setString(1, row[0]); ps.setString(2, row[1]); ps.setString(3, row[2]); ps.addBatch(); } ps.executeBatch(); conn.commit(); // 每批一个事务 return batch.size(); } catch (SQLException e) { // 处理失败逻辑:可以尝试逐条插入以定位脏数据 return 0; } }这里我有几个刻意为之的设计:
第一,每个批次一个事务。有人为了追求极致性能,把几万行放进一个事务,最后如果失败了全部回滚,代价非常大。每批一个事务,失败时只回滚当前批,影响面可控。
第二,try-with-resources确保连接一定归还。这个写习惯了没什么,但很多初学者容易在异常路径上忘记归还连接,最后把连接池耗尽。
第三,批量大小不是越大越好。我实测过,在MySQL里单批500到2000行是最优区间。超过5000行时,单条多值SQL过大,网络包分片反而变慢,数据库解析复杂度也上升。
3.5 脏数据定位与重试策略
并行导入最大的痛点之一就是:某一行格式有问题,整批失败,你却不知道是具体哪一行。
我的做法是这样的:
} catch (SQLException e) { // 批量执行失败时,逐条执行,找出脏数据 try (Connection conn = dataSource.getConnection()) { for (String[] row : batch) { try (PreparedStatement ps = conn.prepareStatement(sql)) { ps.setString(1, row[0]); ps.setString(2, row[1]); ps.setString(3, row[2]); ps.execute(); } catch (SQLException ex) { log.error("脏数据行异常, 内容: {},错误: {}", String.join(",", row), ex.getMessage()); } } conn.commit(); } }这段逻辑虽然损失一点性能,但只会在批量写入失败时触发,不影响正常流程的吞吐。实际项目里,我还习惯在每条原始解析后给数据加上行号,这样定位脏数据时能直接告诉用户"CSV第1382行有问题",体验完全不一样。
4. 参数调优与性能实测
4.1 影响导入速度的四个关键参数
并行方案有效的前提是各项参数配到位,否则效果会大打折扣。
第一,JDBC URL 参数。MySQL的JDBC URL除了上面提到的rewriteBatchedStatements=true,我通常还会加这样几个:
jdbc:mysql://host:3306/db?useUnicode=true&characterEncoding=utf8&rewriteBatchedStatements=true&useServerPrepStmts=true&useCompression=trueuseCompression=true在数据量大、网络带宽有限的场景下效果非常明显,CPU换带宽,多数情况下是划算的。
第二,连接池核心参数。HikariCP中有三四个参数决定导入时的表现:
maximumPoolSize:最大连接数,建议等于线程数。minimumIdle:最小空闲连接数,导入场景设1就行,没必要维护一堆空闲连接。connectionTimeout:取连接的超时时间,导入高峰期线程多,要设长一点,比如30秒。maxLifetime:连接最大存活时间,注意要小于数据库wait_timeout。
第三,批大小。前面提过500到2000行比较合适,我建议固定用1000行起步,实测下来比较均衡。
第四,数据库端配置。如果是自建的MySQL,注意这几个参数:
# my.cnf 中建议关注 max_allowed_packet = 64M innodb_buffer_pool_size >= 物理内存的60% innodb_flush_log_at_trx_commit = 2 # 导入场景下,可以接受性能优先 sync_binlog = 0 或 Ninnodb_flush_log_at_trx_commit=1是最安全但最慢的模式,每条事务提交都要刷盘。导入任务可以降到2甚至0,但要确保能接受极端情况下的少量数据丢失。生产环境建议先确认一下再改,别盲目照抄。
4.2 实测数据:单连接 vs 多连接池差距有多大
我自己在8核16G的云服务器上,MySQL也部署在同一台机器,对一张10个字段的宽表做了测试。表里没有复杂索引,只有主键。数据量是100万行CSV,约500MB。
| 方案 | 参数配置 | 耗时 | 备注 |
|---|---|---|---|
| 单连接逐条插入 | 无优化 | 约40分钟 | 每次INSERT单独提交 |
| 单连接批量插入 | batchSize=1000 | 约2分20秒 | 开了rewriteBatchedStatements |
| 4线程 + 4连接池 | batchSize=1000 | 约55秒 | 线程池4,连接池4 |
| 8线程 + 8连接池 | batchSize=1000 | 约38秒 | 线程池8,连接池8 |
| 16线程 + 16连接池 | batchSize=1000 | 约41秒 | 线程过多,锁等待增加 |
有意思的是,8线程并不是最高的性能拐点,16线程反而略有下降。原因很典型:线程和连接数一多,MySQL的锁竞争和redo日志写入成为新瓶颈。CPU核数是8,线程数超过2倍核数后,上下文切换的开销开始抵消并行收益。
所以要强调一件事:不是线程越多越快,要找到自己环境的拐点。我通常是按2倍核数起步,然后逐步加线程,观察数据库的Threads_running和Innodb_row_lock_current_waits指标,一旦锁等待明显上涨,就说明并发已经到头了。
4.3 导入任务对在线业务的影响控制
并行导入做起来之后,很快会面临第二个问题:导入爽了,线上业务卡了。
这里我建议几个实用手段:
- 限流:在代码里给导入任务加一个总流量控制,比如每秒最多写入多少行,压住瞬时冲击。
- 错峰:大导入尽量安排在业务低峰期。
- 分级:如果数据库支持,把导入任务分配到从库(先写从库,再同步到主库)或者只读副本上,完全不碰主库。
另外,导入期间监控Threads_connected和CPU使用率。如果连接数逼近max_connections,优先降低线程数。
5. 常见问题与排查实录
5.1 导入后出现重复数据
这是高频问题。排查后发现,大部分重复不是SQL写重复了,而是失败重试机制没有做幂等。
比如某个批次执行超时,代码判定为失败并重试,实际数据库那边已经提交成功了,重试就导致同一批数据插了两遍。
解决办法是给目标表加业务唯一索引,导入前用INSERT IGNORE或ON DUPLICATE KEY UPDATE做兜底。更严谨的做法是,CSV导入场景定义好每批次唯一键,比如批次号 + 行号,确保同一批数据只允许出现一次。
5.2 内存溢出
很多人的第一版导入代码是Files.readAllLines()一口气把所有内容load进内存,500MB的文件直接变成2GB的String数组,不OOM才怪。
我后来总结了一个标准姿势:
- 使用
BufferedReader流式读,保证JVM堆里最多只保留一个批次的行。 BlockingQueue用有界队列,容量控制在线程数 × 2的批次数量。- 每批次处理完立即释放引用。
这套组合拳打下来,即使文件有几个GB,Java堆峰值也能控制在几百MB以内。
5.3 导入过程中连接池被耗尽
表现是后台日志疯狂报Connection is not available, request timed out。
这个时候先别急着调大连接池。要分清是连接池真的不够,还是连接泄漏了。
简单的排查方式:在HikariCP配置里打开泄漏检测:
config.setLeakDetectionThreshold(60_000);如果连接从池里拿出去超过60秒没归还,日志会直接打印出获取连接的堆栈。我见过不少"连接池耗尽"最终查出来是某个PreparedStatement没关闭导致连接无法归还的问题。
5.4 CSV内容引起的脏数据问题
这一块最琐碎,但也最影响导入成功率:
- BOM头:Windows记事本保存的CSV带UTF-8 BOM,第一列会多一个不可见字符
\uFEFF。读文件时要用Reader显式处理,或者首行首列replace("\uFEFF", "")。 - Excel导出的CSV:换行符可能是
\r\n,字段内容本身也可能包含逗号和换行。这就是为什么我不建议单纯用split(","),推荐自己写一个状态机解析器或者用成熟的库(如 commons-csv、uniVocity)。 - 空值和null:CSV的空字符串和数据库的NULL不是一回事。导入前要明确规则:空列是写入空字符串还是NULL。
- 字段长度超限:批量插入时整批失败,定位麻烦。解决办法是在解析阶段就做基础校验,长度超过表字段定义的直接标记为脏数据。
5.5 大事务回滚太慢
当batchSize设得太大,或者一批数据里恰好某一行出问题导致整批回滚时,回滚开销会非常可观。尤其是InnoDB,大事务回滚可能比正常提交还慢。
我建议控制在1000行左右一批,这样即使回滚也就几秒钟的事。如果你确实需要大批量一次性导入,可以考虑分批提交,但每批之间用事务边界隔开,确保互不影响。
6. 从"能用"到"好用":稳健性与扩展设计
6.1 多数据源场景下的连接池规划
如果你的导入任务需要同时把同一份CSV写入多个数据库(比如一张表同时落到统计分析库和业务库),多连接池的价值就更明显了。
我的做法是:为每个目标数据库创建一个独立的HikariDataSource,每个数据源独立配置连接池大小。线程池可以共用,但写不同数据库分支的任务建议用独立线程,防止A库慢拖垮B库。
Map<String, HikariDataSource> dataSourceMap = new HashMap<>(); dataSourceMap.put("business", buildImportDataSource(urlA, userA, passA, 8)); dataSourceMap.put("analytics", buildImportDataSource(urlB, userB, passB, 4));这里有个细节值得注意:不同库的写入能力可能差异很大,比如业务库是主库,写入能力明显强于做分析的从库。各自的连接池大小需要独立调优,反向拖累整体任务进度。
6.2 从CSV延伸到Excel、JSON等格式
CSV搞定了,其他格式也好说。原理一样,只是解析器不一样:
- Excel文件用 EasyExcel 或 POI 的流式读取模式,千万别用
WorkbookFactory.create()一次性加载,几万行就会OOM。 - JSON文件用 Jackson 的流式API
JsonParser,边读边解析字段,效果等价于CSV的BufferedReader方案。 - 如果你的CSV表头特别复杂、有多级表头或动态列,可以在解析阶段先维护一个"列名 → 目标字段"的映射关系,把文件中的列顺序和数据库列解耦。开发中遇到的复杂表头Excel需求,本质上就是先做一层表头映射再套用并行写入框架。
6.3 导入任务做得更稳的一个小技巧
最后分享一个我实战中觉得特别有用的设计:给每条插入行增加一个批次ID字段。
导入前生成一个唯一的批次号,这次导入的所有数据都带上这个批次号。万一导入过程出问题需要重导,直接DELETE FROM target_table WHERE batch_id = ?,干净利落。
这个技巧看起来很小,实际用起来救命。有一次我在生产环境跑一个2000万行的导入,跑到后半段发现源文件有数据错误需要重来,多亏有批次ID,几秒钟清掉重导,不然光是清洗脏数据就得折腾半天。
从单连接逐条插入,到多线程多连接池并行导入,这个演进并不复杂,核心思路就三句话:批量提交让单条SQL多干活,多线程让多个SQL同时干活,有界队列让读文件和写数据库的解耦。把这三件事做好,CSV导入的耗时可以从小时级压缩到分钟级。实际操作的时候,多留意我说的那些坑——连接池泄漏、批大小拐点、脏数据定位、批次ID兜底——这些才是真正决定方案能不能平稳落地的关键。