Skip to content

百万行数据导入:流式读取、批量写入与失败处理 ​

运营上传一个 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,中断后从头再跑也不会产生重复数据。

二、整体流程 ​

上传文件xlsx / csv流式读取每次 1000 行校验与转换格式 · 必填有界队列满了就阻塞写入线程 1写入线程 NMySQL批量 INSERT整批失败降级逐行写入定位坏行错误报告:行号 + 原因其余行写入成功
图 1 · 读取端按批次流式解析,经有界队列交给多个写入线程;某一批失败时降级为逐行写入,把坏数据记入错误报告
  1. 读取线程流式解析文件,每 1000 行组成一批。
  2. 校验与转换:格式、必填、取值范围、字典映射,不合法的行直接记入错误报告,不进入队列。
  3. 有界队列连接读取和写入:写入跟不上时队列变满,读取线程被阻塞,内存不会无限增长。
  4. 多个写入线程从队列取批次,用批处理写入数据库,每批一个事务。
  5. 批次失败时降级为逐行写入,定位出错的行。
  6. 导入结束生成结果:成功行数、失败行数、每个失败行的行号与原因。

三、读取:别把整个文件读进内存 ​

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,同样的调用方式在配套实验中编译运行过):

java
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 测试方法 ​

表结构带一个自增主键和一个业务唯一索引,更接近真实业务表:

sql
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 结果 ​

逐行插入 · 自动提交1,918 行/秒批处理 1000 · 未开启重写20,454 行/秒批处理 1000 · 开启重写129,786 行/秒批处理 5000 · 开启重写152,323 行/秒MyBatis foreach · 每批 100079,365 行/秒4 线程 · 批处理 1000 · 重写203,252 行/秒LOAD DATA LOCAL INFILE322,061 行/秒MySQL 8.4.11 默认配置 · Connector/J 8.0.27 · MyBatis 3.5.19 · 表含主键与一个唯一索引 · 3 次采样中位数
图 2 · 同样写入 20 万行,逐行插入每秒不到 2000 行;开启 rewriteBatchedStatements 后快 68 倍,加 4 个线程再快约一半,LOAD DATA 仍然最快;MyBatis 的 foreach 多行插入介于未重写与重写之间(横轴为对数刻度)
方式耗时行/秒相对逐行插入
逐行插入,自动提交(2 万行)10.43s1,9181×
批处理 1000 行一批,未开启重写9.78s20,45411×
批处理 1000 行一批,开启重写1.54s129,78668×
批处理 5000 行一批,开启重写1.31s152,32379×
MyBatis <foreach> 多行插入,1000 行一批,未开启重写2.52s79,36541×
4 个线程,每线程 1000 行一批,开启重写0.98s203,252106×
LOAD DATA LOCAL INFILE0.62s322,061168×

几个观察:

  • 逐行插入慢在每行一次网络往返加一次提交刷盘。 按这个速度,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 写入代码 ​

java
// 连接参数: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 降级逐行写入 ​

可靠的做法是:批次失败就回滚这一批,再逐行重试。

java
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 把逐条插入变成多行插入,二者之间用有界队列形成反压。真正需要花心思的是失败处理:批处理的失败语义会随驱动参数变化,所以要用「批次失败降级逐行」来定位坏行,用业务唯一键让整个导入可以安全重跑。


配套实验

参考资料

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权