核心是分片异步入库:将大任务按100–500条切片,用自定义io线程池(非forkjoinpool)并行执行每批事务性插入,再通过completablefuture.allof汇总结果并捕获executionexception处理异常,同时确保数据库连接池与并发数匹配、事务边界清晰。

用 CompletableFuture 并行分片批量入库,核心是:把大任务拆成多个小批次 → 每批用独立线程异步执行插入 → 汇总结果或处理异常。关键不是“多线程”本身,而是合理控制并发数、避免数据库连接打满、保证事务边界清晰。
分片逻辑:按数量切分 List
假设你有一万个实体对象要入库,不建议单次 insert 一万个(可能超 MySQL packet limit 或锁表太久),通常每批 100–500 条较稳妥。用 Java 8 的 Stream 或传统 for 循环切片:
public static <t> List<list>> partition(List<t> list, int batchSize) {
List<list>> partitions = new ArrayList();
for (int i = 0; i
</list></t></list></t>
调用示例:List<list>> batches = partition(users, 200);</list>
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
并行提交:用 supplyAsync + 自定义线程池
别用默认 ForkJoinPool(它适合 CPU 密集型,而 DB 写入是 I/O 密集型)。创建固定大小的线程池,比如 4–8 个线程(取决于数据库连接池大小和机器核数):
ExecutorService dbExecutor = Executors.newFixedThreadPool(6);
List<completablefuture>> futures = batches.stream()
.map(batch -> CompletableFuture.supplyAsync(() -> {
// 每批在一个事务内执行批量插入
return userMapper.insertBatch(batch); // 返回影响行数
}, dbExecutor))
.collect(Collectors.toList());
</completablefuture>
- 每个
supplyAsync对应一个批次,异步提交到线程池 -
insertBatch方法内部应使用 MyBatis 的<foreach></foreach>或 JdbcTemplate.batchUpdate,确保单批原子性 - 不要在 lambda 里开启新事务(如 @Transactional),事务需由该方法自身管理
等待完成与错误处理
用 CompletableFuture.allOf 等待全部结束,但注意它不返回结果也不抛异常;推荐用 thenCollect 风格链式收集,或手动遍历 future 获取结果:
// 收集所有结果(含异常)
List<integer> results = futures.stream()
.map(future -> {
try {
return future.get(); // 阻塞获取,生产环境建议加超时:future.get(30, TimeUnit.SECONDS)
} catch (ExecutionException e) {
Throwable cause = e.getCause();
log.error("批量插入失败", cause);
throw new RuntimeException("入库失败", cause);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
})
.collect(Collectors.toList());
int total = results.stream().mapToInt(Integer::intValue).sum();
log.info("共写入 {} 条记录", total);
</integer>
- 务必捕获
ExecutionException,它包装了业务异常(如唯一键冲突、空指针) - 避免无超时的
get(),防止某个批次卡死拖垮整个流程 - 如果某批失败需整体回滚,应在最外层 catch 后手动清理已入库数据(补偿事务),或改用分布式事务框架
进阶提醒:连接池与事务隔离
MyBatis 默认每个 SqlSession 使用独立数据库连接,只要你的 insertBatch 在同一个 mapper 接口方法里执行,就天然在单个连接+事务中完成。但要注意:
- HikariCP 连接池最大连接数 ≥ 并发线程数,否则线程会阻塞在 getConnection()
- MySQL 默认事务隔离级别是 REPEATABLE READ,大批量插入期间其他查询可能被锁或变慢,可考虑临时降级为 READ COMMITTED(需评估一致性要求)
- 插入前做必要去重或幂等校验(如先 select for update 或用 insert ignore / on duplicate key update)
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南










