百万行数据导入:流式读取、批量写入与失败处理
运营上传一个 100 万行的 Excel,服务先是 OOM,改成分批后又慢得超时,好不容易跑完,又发现其中几百行格式错误的数据让整批都没写进去。导入功能的三个难点是内存、速度和错误处理,而且三者的解法互相牵制。
本文先用 JDBC 实测几种写入方式的速度差异,再看一个容易被忽略的细节:开启 rewriteBatchedStatements 后,批处理失败时的行为会完全改变,这直接决定了错误处理该怎么设计。
一、先说结论
- 内存问题靠流式读取解决:按行解析、每攒够一批就交出去,内存占用只和批大小有关,与文件大小无关。
- 速度的关键是 JDBC 连接参数
rewriteBatchedStatements=true。 实测写入 20 万行:逐行插入每秒 1,918 行,普通批处理 20,454 行,开启重写后 129,786 行,再加 4 个线程 203,252 行;LOAD DATA为 322,061 行,仍然最快。 - 开启重写后,一行出错会让整批失败。 实测一批 6 行中有 1 行主键冲突:不开启重写时其余 5 行写入成功;开启后一行都没有写入。
- 错误处理要分两层:批次失败后降级为逐行写入,找出坏行记入报告,其余行正常写入;数据校验尽量在写入之前完成。
- 导入要能重跑:用业务唯一键加
INSERT ... ON DUPLICATE KEY UPDATE或INSERT IGNORE,中断后从头再跑也不会产生重复数据。
二、整体流程
- 读取线程流式解析文件,每 1000 行组成一批。
- 校验与转换:格式、必填、取值范围、字典映射,不合法的行直接记入错误报告,不进入队列。
- 有界队列连接读取和写入:写入跟不上时队列变满,读取线程被阻塞,内存不会无限增长。
- 多个写入线程从队列取批次,用批处理写入数据库,每批一个事务。
- 批次失败时降级为逐行写入,定位出错的行。
- 导入结束生成结果:成功行数、失败行数、每个失败行的行号与原因。
三、读取:别把整个文件读进内存
Apache POI 的 XSSFWorkbook 会把整个工作簿解析成对象树。实测一个 100 万行、10 列、约 63 MB 的 xlsx:XSSFWorkbook 在 6 GB 堆下抛出 OutOfMemoryError,给到 8 GB 才读完;同一个文件用下面的 Fesod 按批读取,256 MB 堆就够了,堆使用峰值约 159 MB。流式读取只在内存中保留当前这一批。
常用的做法:
- Excel:EasyExcel 是国内使用最广的流式读取库,已停止维护;原作者团队后续的 FastExcel 在 2025 年 9 月进入 Apache 孵化器,改名为 Apache Fesod。新项目可以直接使用 Fesod,入口类是
FesodSheet。 - CSV:能让用户上传 CSV 就优先 CSV,用
BufferedReader逐行读取即可,解析成本远低于 Excel。
使用 Fesod 按批读取的写法如下(代码片段,依赖 org.apache.fesod:fesod-sheet:2.0.2-incubating,同样的调用方式在配套实验中编译运行过):
BlockingQueue<List<OrderRow>> queue = new ArrayBlockingQueue<>(16); // 最多积压 16 批
FesodSheet.read(file, OrderRow.class, new PageReadListener<OrderRow>(batch -> {
try {
queue.put(validate(batch)); // 队列满时阻塞读取线程
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException(e);
}
}, 1000)).sheet().doRead();常见建议中的「多线程并发读取多个 Sheet」只适用于数据本来就分布在多个 Sheet 的文件。单个 Sheet 的解析本身是顺序的,瓶颈通常也不在解析,而在写入。
四、写入:几种方式实测
4.1 测试方法
表结构带一个自增主键和一个业务唯一索引,更接近真实业务表:
CREATE TABLE import_order (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
order_no BIGINT NOT NULL,
user_id BIGINT NOT NULL,
sku VARCHAR(32) NOT NULL,
qty INT NOT NULL,
remark VARCHAR(64) NOT NULL,
UNIQUE KEY uk_order_no (order_no)
);每种方式重新建表后写入 20 万行(逐行插入太慢,只写 2 万行),采样 3 次取中位数。环境是 JDK 21、MySQL Connector/J 8.0.27、默认配置的 MySQL 8.4.11(开启 binlog,innodb_flush_log_at_trx_commit 与 sync_binlog 都为 1)。
4.2 结果
| 方式 | 耗时 | 行/秒 | 相对逐行插入 |
|---|---|---|---|
| 逐行插入,自动提交(2 万行) | 10.43s | 1,918 | 1× |
| 批处理 1000 行一批,未开启重写 | 9.78s | 20,454 | 11× |
| 批处理 1000 行一批,开启重写 | 1.54s | 129,786 | 68× |
| 批处理 5000 行一批,开启重写 | 1.31s | 152,323 | 79× |
MyBatis <foreach> 多行插入,1000 行一批,未开启重写 | 2.52s | 79,365 | 41× |
| 4 个线程,每线程 1000 行一批,开启重写 | 0.98s | 203,252 | 106× |
LOAD DATA LOCAL INFILE | 0.62s | 322,061 | 168× |
几个观察:
- 逐行插入慢在每行一次网络往返加一次提交刷盘。 按这个速度,100 万行需要 8 分半。
- 只用批处理、不开启重写,只省掉了提交次数。 驱动仍然逐条发送
INSERT,网络往返次数没变。 - 开启
rewriteBatchedStatements=true后,驱动把一批INSERT改写成一条多行插入INSERT ... VALUES (...), (...), ...,一次往返写入一整批。这是收益最大的一步。 - 批大小从 1000 加到 5000,只快了 17%,但一批失败时要重试的数据变成 5 倍,1000 左右是比较均衡的选择。改写后的单条语句不能超过
max_allowed_packet(MySQL 8.4 默认 64MB),行很宽时要相应减小批大小。 - 多线程并行在这台机器上快了约 50%(容器限制 2 CPU)。线程数不是越多越好:数据库的写入能力、唯一索引的检查、自增锁都会成为瓶颈,而且导入会挤占线上业务的数据库资源,要限制并发并安排在低峰期。
LOAD DATA 比 4 线程的批处理还快约 60%,但需要服务端开启 local_infile、客户端开启 allowLoadLocalInfile,出于安全考虑很多生产环境会禁用。在应用内导入时,开启重写的批处理通常已经足够。
4.3 写入代码
// 连接参数:jdbc:mysql://host:3306/db?rewriteBatchedStatements=true
void writeBatch(DataSource ds, List<OrderRow> batch) throws SQLException {
String sql = "INSERT INTO import_order (order_no, user_id, sku, qty, remark) VALUES (?, ?, ?, ?, ?)";
try (Connection c = ds.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
c.setAutoCommit(false);
try {
for (OrderRow r : batch) {
ps.setLong(1, r.orderNo());
ps.setLong(2, r.userId());
ps.setString(3, r.sku());
ps.setInt(4, r.qty());
ps.setString(5, r.remark());
ps.addBatch();
}
ps.executeBatch();
c.commit();
} catch (SQLException e) {
c.rollback();
throw e;
}
}
}使用 MyBatis 时,常见的做法是用 <foreach> 拼接多行 VALUES。实测每批 1000 行、连接不开启重写,每秒 79,365 行:是未重写批处理的 4 倍,但只有开启重写的批处理的六成。用 <foreach> 还要自己控制每批的行数,避免 SQL 超过 max_allowed_packet;能改连接参数时,JDBC 批处理加 rewriteBatchedStatements=true 更快;MyBatis 的 ExecutorType.BATCH 底层走的也是 JDBC 批处理。
五、失败处理:批次失败时发生了什么
5.1 实测:一批里有一行重复
一批 6 行,第 4 行的 order_no 与第 2 行重复,捕获异常后故意提交,观察哪些行被写入:
| 设置 | getUpdateCounts() | 提交后表中的数据 |
|---|---|---|
| 未开启重写 | [1, 1, 1, -3, 1, 1] | 1, 2, 3, 5, 6 |
| 开启重写 | [-3, -3, -3, -3, -3, -3] | 空 |
-3 是 Statement.EXECUTE_FAILED。
- 未开启重写时,驱动逐条执行,出错的那条失败,其余照常执行(Connector/J 的
continueBatchOnError默认为true)。如果捕获异常后提交,就会得到一批「部分成功」的数据。 - 开启重写后,整批是一条语句,一行出错整条语句失败,这一批一行都没有写入。
两种行为都不能直接用来「跳过坏行继续导入」:前者的结果依赖驱动的执行方式,后者整批丢失。
5.2 降级逐行写入
可靠的做法是:批次失败就回滚这一批,再逐行重试。
void writeWithFallback(DataSource ds, List<OrderRow> batch, ImportReport report) {
try {
writeBatch(ds, batch);
report.success(batch.size());
} catch (SQLException batchError) {
for (OrderRow row : batch) {
try {
writeBatch(ds, List.of(row));
report.success(1);
} catch (SQLException rowError) {
report.fail(row.lineNo(), rowError.getMessage());
}
}
}
}降级只发生在出错的批次上,正常批次仍然走批处理。坏数据很多时降级会很慢,所以要把能提前发现的问题放到校验阶段:格式、必填、长度、枚举值、文件内部的重复键。
5.3 让导入可以重跑
导入可能在任何位置中断:应用发布、超时、数据库切换。设计成可重跑,比实现精确的断点续传简单得多:
- 用业务唯一键(如订单号)建唯一索引;
- 按需求选择冲突时的行为:
INSERT IGNORE跳过已存在的行,或INSERT ... ON DUPLICATE KEY UPDATE覆盖为文件中的值; - 记录导入任务的状态和已处理的批次号,重跑时可以跳过已确认完成的批次,作为优化而不是正确性的保证。
注意 INSERT IGNORE 会把一些其他错误也降级为警告,例如数据被截断。对数据质量要求严格时,用 ON DUPLICATE KEY UPDATE,并在校验阶段拦住格式问题。
5.4 关于「整个导入放在一个事务里」
常见的「用 @Transactional 保证失败时整体回滚」,对百万行导入不适用:一个事务写 100 万行会产生巨大的 Undo 和 Binlog,锁持有时间长,失败回滚同样耗时,replica 也会出现明显的复制延迟,原因与 千万级大表怎么清理数据 中的大删除相同。
如果业务确实要求「全部成功才生效」,可以先导入到临时表或带批次号的暂存表,校验通过后再用一条语句切换状态,或者分批从暂存表迁移到正式表。
六、上线前的检查清单
- [ ] 读取使用流式 API,内存占用与文件大小无关
- [ ] 读写之间有有界队列,写入慢时读取会被阻塞
- [ ] 连接参数开启
rewriteBatchedStatements=true,批大小在 500—2000 之间 - [ ] 写入线程数有上限,导入任务限制在业务低峰期
- [ ] 批次失败时降级逐行写入,错误报告带行号和原因
- [ ] 业务唯一键有唯一索引,导入可以安全重跑
- [ ] 导入是异步任务,前端轮询进度,不占用 HTTP 请求线程等待
- [ ] 文件大小和行数有上限,防止恶意或误操作上传超大文件
七、常见误区
- 「用了
addBatch就是批量插入」:不开启rewriteBatchedStatements时,驱动仍然逐条发送语句。 - 「批越大越快」:实测从 1000 到 5000 只快 17%,失败时的重试代价却变成 5 倍。
- 「整个导入放一个事务最安全」:百万行的大事务会带来锁、Undo、Binlog 和复制延迟问题。
- 「批处理失败后提交,就能保留正确的行」:开启重写时整批都没有写入,未开启时结果取决于驱动的执行方式。
- 「并发线程越多导入越快」:数据库写入能力有上限,还会影响线上业务。
小结
导入百万行数据,读取端靠流式解析控制内存,写入端靠 rewriteBatchedStatements 把逐条插入变成多行插入,二者之间用有界队列形成反压。真正需要花心思的是失败处理:批处理的失败语义会随驱动参数变化,所以要用「批次失败降级逐行」来定位坏行,用业务唯一键让整个导入可以安全重跑。
配套实验
- codesphere-labs/storage/mysql-bulk-import:七种写入方式(含 MyBatis
<foreach>)各采样 3 次,以及一批 6 行中有 1 行主键冲突时两种驱动设置下的失败语义(验证记录) - codesphere-labs/storage/excel-import-memory:100 万行 xlsx 用 POI
XSSFWorkbook与 Apache Fesod 读取的堆内存(验证记录)
参考资料