ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Spring Boot海量数据高效写入实战:MyBatis-Plus批处理与多线程优化

Spring Boot海量数据高效写入实战:MyBatis-Plus批处理与多线程优化 最近在开发一个需要处理大量数据导入和实时更新的项目时遇到了一个棘手的问题如何高效、安全地将外部数据源如Excel、CSV的数据同步到业务数据库并确保在数据量激增时系统依然稳定。传统的逐条插入或简单批处理在面对十万、百万级数据时往往力不从心容易导致数据库连接池耗尽、事务超时甚至引发线上服务雪崩。本文将围绕“海量数据高效写入”这一核心主题深入拆解从数据读取、分批处理、事务控制到性能监控的全流程实战方案。不同于零散的代码片段我会结合一个模拟的“用户信息批量导入”场景提供一套完整、可复现的Java/Spring Boot解决方案涵盖MyBatis-Plus、线程池、数据库连接池等关键组件的深度应用。无论你是正在优化现有数据同步功能还是为即将到来的大数据量场景做准备这篇文章都能提供从理论到落地的闭环指导。1. 背景与核心概念为什么需要专门的数据写入策略在业务开发中数据写入是最基础的操作之一。但当数据量从“几十条”变为“几十万条”时简单的for循环插入就会暴露出诸多问题性能瓶颈网络I/O和数据库I/O次数与数据量成正比耗时线性增长用户体验差。资源耗尽每条插入都占用一个数据库连接海量数据瞬间打满连接池导致其他正常请求无法获取连接。事务过大将所有插入放在一个大事务中事务日志暴涨锁持有时间过长极易造成死锁或主从延迟。容错性差中间某条数据失败可能导致整个批次回滚或者需要复杂的重试和补偿机制。因此我们需要一套**批处理Batch Processing**策略。其核心思想是化整为零分而治之。分批Batch将海量数据分割成多个大小适宜的批次例如每1000条一批。批操作Batch Operation每个批次的数据通过一次数据库交互完成写入大幅减少网络和I/O开销。可控的事务每个批次或几个批次作为一个独立事务避免大事务便于失败重试。本文将实现的方案正是基于这种思想并结合了连接池、线程池等技术构建一个健壮的海量数据写入服务。2. 环境准备与版本说明为了完整演示我们需要搭建一个标准的Spring Boot项目环境。以下是本文示例所使用的主要组件及版本你可以根据实际项目情况进行调整。核心环境JDK: 1.8 或 11 (推荐11本文示例基于11)Spring Boot: 2.7.x (本文使用2.7.18)构建工具: Maven 3.6主要依赖 (pom.xml):dependencies !-- Spring Boot Starter -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-jdbc/artifactId /dependency !-- MyBatis-Plus (极大简化数据库操作) -- dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.3.1/version /dependency !-- 数据库驱动 (以MySQL为例) -- dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId scoperuntime/scope version8.0.33/version /dependency !-- 连接池 (Spring Boot 2.7默认使用HikariCP这里显式声明) -- dependency groupIdcom.zaxxer/groupId artifactIdHikariCP/artifactId /dependency !-- 工具类 -- dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId optionaltrue/optional /dependency dependency groupIdorg.apache.commons/groupId artifactIdcommons-collections4/artifactId version4.4/version /dependency dependency groupIdorg.apache.commons/groupId artifactIdcommons-lang3/artifactId /dependency !-- 测试 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency /dependencies数据库表结构示例我们以一个简单的user_info表为例模拟用户信息导入。CREATE TABLE user_info ( id bigint(20) NOT NULL AUTO_INCREMENT COMMENT 主键ID, username varchar(64) NOT NULL COMMENT 用户名, email varchar(128) DEFAULT NULL COMMENT 邮箱, age int(11) DEFAULT NULL COMMENT 年龄, create_time datetime DEFAULT CURRENT_TIMESTAMP COMMENT 创建时间, update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT 更新时间, PRIMARY KEY (id), UNIQUE KEY uk_username (username) COMMENT 用户名唯一索引 ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT用户信息表;项目结构预览src/main/java/com/example/batchimport/ ├── BatchImportApplication.java // 启动类 ├── config/ │ ├── ThreadPoolConfig.java // 线程池配置 │ └── DataSourceConfig.java // 数据源配置可选 ├── entity/ │ └── UserInfo.java // 实体类 ├── mapper/ │ └── UserInfoMapper.java // MyBatis-Plus Mapper接口 ├── service/ │ ├── IUserService.java // 服务接口 │ └── impl/ │ └── UserServiceImpl.java // 服务实现类核心逻辑在此 └── controller/ └── DataImportController.java // 提供导入接口3. 核心组件与原理拆解在实现批处理导入前需要理解几个关键组件的配置与原理它们是保障方案稳定性的基石。3.1 数据库连接池 (HikariCP) 关键配置连接池管理着与数据库的连接其配置直接影响批处理的并发能力和稳定性。Spring Boot默认使用HikariCP在application.yml中配置如下spring: datasource: driver-class-name: com.mysql.cj.jdbc.Driver url: jdbc:mysql://localhost:3306/test_db?useUnicodetruecharacterEncodingutf8useSSLfalseserverTimezoneAsia/ShanghairewriteBatchedStatementstrue # 关键参数 username: root password: yourpassword hikari: # 连接池名称 pool-name: BatchImportHikariPool # 最小空闲连接数 minimum-idle: 10 # 最大连接池大小 (必须根据数据库和机器性能调整) maximum-pool-size: 50 # 连接最大存活时间毫秒默认30分钟防止数据库端连接僵死 max-lifetime: 1800000 # 连接超时时间毫秒 connection-timeout: 30000 # 空闲连接超时时间毫秒 idle-timeout: 600000 # 连接测试查询 connection-test-query: SELECT 1关键参数解释rewriteBatchedStatementstrue:这是MySQL批处理性能提升的关键。它允许JDBC驱动将多个INSERT语句重写为单个多值INSERT语句如INSERT INTO table VALUES (...), (...), (...)大幅减少网络往返。maximum-pool-size: 最大连接数。此值需谨慎设置过小会导致并发批处理任务等待连接过大则可能压垮数据库。一般建议设置为(核心数 * 2) 有效磁盘数的估算值并结合压测调整。minimum-idle: 最小空闲连接。保持一定数量的“热”连接避免突发请求时创建连接的延迟。3.2 MyBatis-Plus 的批处理支持MyBatis-Plus的SqlInjector和ServiceImpl类提供了开箱即用的批处理方法saveBatch。但其默认实现可能不是最优的。我们需要关注其底层是如何工作的。默认的saveBatch逻辑遍历实体集合。为每个实体生成一条INSERTSQL。使用SqlSession的executeBatch方法执行。该方法依赖于JDBC的addBatch和executeBatch机制。配合MySQL的rewriteBatchedStatementstrue驱动会将多个INSERT合并。性能瓶颈即使合并了SQL如果一次传入10万条数据saveBatch默认仍会尝试在一个数据库会话中处理可能造成内存溢出或事务过大。因此我们需要手动控制批次大小。3.3 线程池 (ThreadPoolTaskExecutor) 配置为了进一步提升导入速度我们可以利用多线程并行处理多个批次。Spring提供的ThreadPoolTaskExecutor是一个很好的选择。import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.ThreadPoolExecutor; Configuration public class ThreadPoolConfig { Bean(batchImportThreadPool) public ThreadPoolTaskExecutor batchImportThreadPool() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); // 核心线程数服务器核心数 executor.setCorePoolSize(Runtime.getRuntime().availableProcessors()); // 最大线程数核心数 * 2 (根据I/O密集型任务调整) executor.setMaxPoolSize(Runtime.getRuntime().availableProcessors() * 2); // 队列容量不宜过大避免内存堆积和响应延迟 executor.setQueueCapacity(100); // 线程名前缀 executor.setThreadNamePrefix(batch-import-); // 拒绝策略CallerRunsPolicy - 由调用者线程通常是HTTP请求线程自己执行任务。 // 这可以防止任务被丢弃但会降低导入接口的响应速度是一种简单的背压机制。 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); // 核心线程超时回收 executor.setAllowCoreThreadTimeOut(true); // 线程空闲存活时间秒 executor.setKeepAliveSeconds(60); executor.initialize(); return executor; } }配置要点队列容量QueueCapacity设置了队列大小。当线程数达到核心数且队列未满时新任务会进入队列等待。如果队列也满了且线程数未达最大数则会创建新线程。队列不宜过大否则任务积压内存且响应延迟高。拒绝策略RejectedExecutionHandler当线程池和队列都满了如何处理新任务CallerRunsPolicy让调用线程自己执行是一种简单的“降级”策略保证任务不丢失但会影响主线程。对于数据导入也可以考虑DiscardPolicy丢弃或记录日志后丢弃具体看业务对数据完整性的要求。4. 完整实战海量用户数据批量导入现在我们将把上述组件组合起来实现一个完整的、支持多线程分批导入的服务。4.1 创建实体类与Mapper首先定义与数据库表对应的实体类和MyBatis-Plus的Mapper接口。实体类UserInfo.java:package com.example.batchimport.entity; import com.baomidou.mybatisplus.annotation.IdType; import com.baomidou.mybatisplus.annotation.TableId; import com.baomidou.mybatisplus.annotation.TableName; import lombok.Data; import java.time.LocalDateTime; Data TableName(user_info) public class UserInfo { TableId(type IdType.AUTO) private Long id; private String username; private String email; private Integer age; private LocalDateTime createTime; private LocalDateTime updateTime; }Mapper接口UserInfoMapper.java:package com.example.batchimport.mapper; import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.example.batchimport.entity.UserInfo; import org.apache.ibatis.annotations.Mapper; Mapper public interface UserInfoMapper extends BaseMapperUserInfo { // 继承BaseMapper已包含基本的CRUD方法 }4.2 服务层核心实现服务层是业务逻辑的核心。我们将实现一个UserServiceImpl它包含单线程分批插入。多线程并行分批插入。带有简单失败重试的插入。服务接口IUserService.java:package com.example.batchimport.service; import com.example.batchimport.entity.UserInfo; import java.util.List; public interface IUserService { /** * 单线程批量保存内部处理分批 * param userList 用户列表 * param batchSize 每批大小 * return 是否全部成功 */ boolean saveBatchSingleThread(ListUserInfo userList, int batchSize); /** * 多线程并行批量保存 * param userList 用户列表 * param batchSize 每批大小 * return 是否全部成功 */ boolean saveBatchMultiThread(ListUserInfo userList, int batchSize); /** * 带重试的批量保存单线程 * param userList 用户列表 * param batchSize 每批大小 * param maxRetries 最大重试次数 * return 是否全部成功 */ boolean saveBatchWithRetry(ListUserInfo userList, int batchSize, int maxRetries); }服务实现UserServiceImpl.java(核心代码较长分块解析):package com.example.batchimport.service.impl; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import com.example.batchimport.entity.UserInfo; import com.example.batchimport.mapper.UserInfoMapper; import com.example.batchimport.service.IUserService; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections4.ListUtils; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.util.StopWatch; import javax.annotation.Resource; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executor; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; Service Slf4j public class UserServiceImpl extends ServiceImplUserInfoMapper, UserInfo implements IUserService { Resource private UserInfoMapper userInfoMapper; // 注入我们配置的线程池 Resource Qualifier(batchImportThreadPool) private Executor batchImportExecutor; // 默认批次大小 private static final int DEFAULT_BATCH_SIZE 1000; /** * 单线程分批插入 */ Override Transactional(rollbackFor Exception.class) // 整个方法在一个事务内 public boolean saveBatchSingleThread(ListUserInfo userList, int batchSize) { if (userList null || userList.isEmpty()) { log.warn(待插入用户列表为空); return true; } if (batchSize 0) { batchSize DEFAULT_BATCH_SIZE; } StopWatch stopWatch new StopWatch(); stopWatch.start(单线程分批导入); // 使用 commons-collections4 的 ListUtils 进行分区 ListListUserInfo partitions ListUtils.partition(userList, batchSize); log.info(数据总量: {} 划分为 {} 批 每批大小: {}, userList.size(), partitions.size(), batchSize); AtomicBoolean successFlag new AtomicBoolean(true); AtomicInteger processedCount new AtomicInteger(0); for (int i 0; i partitions.size(); i) { ListUserInfo batch partitions.get(i); try { // 使用MyBatis-Plus的saveBatch方法传入批次数据。 // 注意此处的saveBatch是父类ServiceImpl提供的它内部会处理批处理。 boolean result super.saveBatch(batch); if (!result) { log.error(第 {} 批数据插入失败 本批数据量: {}, i 1, batch.size()); successFlag.set(false); // 这里可以根据业务决定是继续还是中断。为了演示我们记录错误但继续。 } processedCount.addAndGet(batch.size()); if ((i 1) % 10 0 || i partitions.size() - 1) { log.info(进度: {}/{} 已处理 {} 条记录, i 1, partitions.size(), processedCount.get()); } } catch (Exception e) { log.error(处理第 {} 批数据时发生异常 本批数据量: {}, i 1, batch.size(), e); successFlag.set(false); // 事务注解会在此处回滚整个方法的所有操作。如果需要部分成功需要调整事务边界。 } } stopWatch.stop(); log.info(单线程导入完成总耗时: {} 毫秒 成功率: {}, stopWatch.getTotalTimeMillis(), successFlag.get()); return successFlag.get(); } /** * 多线程并行分批插入 * 注意此方法本身不开启事务事务由每个子任务内部控制。 */ Override public boolean saveBatchMultiThread(ListUserInfo userList, int batchSize) { if (userList null || userList.isEmpty()) { log.warn(待插入用户列表为空); return true; } if (batchSize 0) { batchSize DEFAULT_BATCH_SIZE; } StopWatch stopWatch new StopWatch(); stopWatch.start(多线程并行导入); ListListUserInfo partitions ListUtils.partition(userList, batchSize); log.info(数据总量: {} 划分为 {} 批 每批大小: {}, userList.size(), partitions.size(), batchSize); // 用于收集异步任务 ListCompletableFutureBoolean futures new ArrayList(); AtomicInteger processedBatches new AtomicInteger(0); AtomicInteger totalProcessed new AtomicInteger(0); for (int i 0; i partitions.size(); i) { ListUserInfo batch partitions.get(i); final int batchIndex i; // 为每个批次提交一个异步任务 CompletableFutureBoolean future CompletableFuture.supplyAsync(() - { // 每个批次在自己的事务中执行 return executeBatchInsert(batch, batchIndex, processedBatches, totalProcessed, partitions.size()); }, batchImportExecutor); // 使用我们自定义的线程池 futures.add(future); } // 等待所有异步任务完成 CompletableFutureVoid allFutures CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])); try { allFutures.join(); // 阻塞直到所有任务完成 } catch (Exception e) { log.error(多线程任务执行过程中发生异常, e); } // 检查所有批次的结果 boolean allSuccess futures.stream().allMatch(f - { try { return f.get(); // 获取每个任务的结果 } catch (Exception e) { log.error(获取异步任务结果异常, e); return false; } }); stopWatch.stop(); log.info(多线程导入完成总耗时: {} 毫秒 全部成功: {}, stopWatch.getTotalTimeMillis(), allSuccess); return allSuccess; } /** * 执行单个批次插入的内部方法带事务 */ Transactional(rollbackFor Exception.class, propagation Propagation.REQUIRES_NEW) // 每个批次独立事务 public boolean executeBatchInsert(ListUserInfo batch, int batchIndex, AtomicInteger processedBatches, AtomicInteger totalProcessed, int totalBatches) { try { boolean result super.saveBatch(batch); int currentBatch processedBatches.incrementAndGet(); int currentTotal totalProcessed.addAndGet(batch.size()); if (currentBatch % 5 0 || currentBatch totalBatches) { log.info(线程[{}] 进度: {}/{} 累计已处理 {} 条记录, Thread.currentThread().getName(), currentBatch, totalBatches, currentTotal); } return result; } catch (Exception e) { log.error(线程[{}] 处理第 {} 批数据时失败 本批数据量: {}, Thread.currentThread().getName(), batchIndex 1, batch.size(), e); // 抛出异常触发事务回滚。由于是独立事务只回滚本批次。 throw new RuntimeException(批次插入失败, e); } } /** * 带重试的单线程分批插入 */ Override public boolean saveBatchWithRetry(ListUserInfo userList, int batchSize, int maxRetries) { if (userList null || userList.isEmpty()) { return true; } if (batchSize 0) batchSize DEFAULT_BATCH_SIZE; if (maxRetries 0) maxRetries 3; ListListUserInfo partitions ListUtils.partition(userList, batchSize); AtomicBoolean finalSuccess new AtomicBoolean(true); for (int i 0; i partitions.size(); i) { ListUserInfo batch partitions.get(i); boolean batchSuccess false; int retryCount 0; Exception lastException null; // 重试逻辑 while (retryCount maxRetries !batchSuccess) { try { // 每次重试都使用新的事务 batchSuccess retryBatchInsert(batch); } catch (Exception e) { lastException e; retryCount; log.warn(第 {} 批数据第 {} 次重试失败 等待后重试..., i 1, retryCount, e); try { // 简单的指数退避等待 Thread.sleep((long) (Math.pow(2, retryCount) * 1000)); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } } } if (!batchSuccess) { log.error(第 {} 批数据在 {} 次重试后仍失败 本批数据将被跳过。最后异常: , i 1, maxRetries, lastException); finalSuccess.set(false); // 记录失败批次后续可人工处理或进入死信队列 // recordFailedBatch(batch, lastException); } } return finalSuccess.get(); } Transactional(rollbackFor Exception.class) public boolean retryBatchInsert(ListUserInfo batch) { return super.saveBatch(batch); } }4.3 控制器层提供接口创建一个简单的Controller来触发导入并模拟生成测试数据。package com.example.batchimport.controller; import com.example.batchimport.entity.UserInfo; import com.example.batchimport.service.IUserService; import lombok.extern.slf4j.Slf4j; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; import java.time.LocalDateTime; import java.util.ArrayList; import java.util.List; import java.util.UUID; RestController RequestMapping(/api/import) Slf4j public class DataImportController { Resource private IUserService userService; /** * 模拟生成测试数据并导入 * param total 生成的数据总量 * param batchSize 每批大小 * param mode 导入模式single-单线程multi-多线程retry-带重试 * return 导入结果 */ PostMapping(/users) public String importUsers(RequestParam(defaultValue 10000) int total, RequestParam(defaultValue 1000) int batchSize, RequestParam(defaultValue single) String mode) { log.info(开始生成 {} 条模拟用户数据批次大小: {}, 模式: {}, total, batchSize, mode); ListUserInfo userList generateMockUsers(total); boolean success false; long startTime System.currentTimeMillis(); switch (mode.toLowerCase()) { case single: success userService.saveBatchSingleThread(userList, batchSize); break; case multi: success userService.saveBatchMultiThread(userList, batchSize); break; case retry: success userService.saveBatchWithRetry(userList, batchSize, 3); break; default: return 不支持的导入模式: mode; } long endTime System.currentTimeMillis(); String resultMsg String.format(导入完成 模式: %s 数据量: %d 批次: %d 总耗时: %d ms 结果: %s, mode, total, batchSize, (endTime - startTime), success ? 成功 : 存在失败); log.info(resultMsg); return resultMsg; } /** * 生成模拟用户数据 */ private ListUserInfo generateMockUsers(int count) { ListUserInfo list new ArrayList(count); for (int i 0; i count; i) { UserInfo user new UserInfo(); user.setUsername(user_ UUID.randomUUID().toString().substring(0, 8) _ i); user.setEmail(user.getUsername() example.com); user.setAge(20 (int)(Math.random() * 30)); // 20-50岁 user.setCreateTime(LocalDateTime.now()); user.setUpdateTime(LocalDateTime.now()); list.add(user); } return list; } }4.4 运行与验证启动应用确保MySQL服务已启动并创建好test_db数据库和user_info表。运行Spring Boot应用。调用接口使用Postman、curl或浏览器访问API。单线程导入1万条数据POST http://localhost:8080/api/import/users?total10000batchSize1000modesingle多线程导入10万条数据POST http://localhost:8080/api/import/users?total100000batchSize2000modemulti带重试导入可模拟网络抖动POST http://localhost:8080/api/import/users?total5000batchSize500moderetry观察日志在应用控制台你将看到详细的分批处理日志、进度和耗时。查询数据库直接查询user_info表确认数据已按预期插入。4.5 结果说明与性能对比通过调整total、batchSize和mode参数你可以直观地感受不同策略的性能差异。一般来说单线程模式逻辑简单事务易控制但总耗时最长。多线程模式充分利用多核CPU和数据库连接池吞吐量最高但事务管理复杂每个批次独立需要关注数据库的并发写入压力。批次大小batchSize需要权衡。太小则网络往返开销大太大则单次数据库操作内存占用高可能超出数据库max_allowed_packet等限制。通常建议在500-2000之间进行压测调优。带重试模式增加了鲁棒性适合网络不稳定或数据库偶发故障的场景但会引入延迟。5. 常见问题与排查思路在实际使用中你可能会遇到以下问题问题现象可能原因排查思路与解决方案插入速度非常慢甚至不如逐条插入1. MySQL连接参数未加rewriteBatchedStatementstrue。2. 批次大小设置不合理如设为1。3. 表有大量索引每次插入需更新索引。1. 检查JDBC URL确保已添加rewriteBatchedStatementstrue。2. 调整batchSize如1000并进行测试。3. 考虑在导入前暂时禁用非关键索引导入后再重建。出现Communications link failure或连接超时1. 数据库连接池配置不当connectionTimeout太小。2. 单批次数据量太大操作时间超过数据库wait_timeout。3. 网络不稳定。1. 适当增大connectionTimeout和数据库的wait_timeout。2. 减少batchSize。3. 实现如示例中的重试机制。多线程导入时数据部分丢失1. 某个子线程任务失败但未正确捕获或处理。2. 事务未正确配置导致部分批次未提交。1. 确保每个异步任务CompletableFuture的结果都被检查如示例中的allFutures.join()和结果遍历。2. 确认Transactional(propagation Propagation.REQUIRES_NEW)正确使用确保每个批次独立提交。数据库CPU或IO飙升1. 并发线程数过多超过数据库处理能力。2. 批次过大产生大事务。1. 降低线程池的maxPoolSize。2. 减少batchSize。3. 考虑在业务低峰期执行批量任务。唯一键冲突Duplicate entry模拟数据或真实数据中存在重复的唯一键如username。1. 在业务层保证数据唯一性或使用ON DUPLICATE KEY UPDATE语义。2. 使用insert ignore忽略重复可能丢失数据。3. 在导入前先对数据进行去重。内存溢出OOM一次性在内存中加载了全部海量数据如从一个大文件读取到List。1. 使用流式读取如Apache Commons CSV的CSVParser或使用BufferedReader逐行处理。2. 分片读取和处理不要一次性加载所有数据到内存。6. 最佳实践与工程建议将海量数据写入从“能用”升级到“好用”和“稳定”还需要考虑以下工程实践监控与可观测性为批处理任务添加详细的日志包括开始时间、结束时间、处理总量、成功/失败数量、耗时等。集成Micrometer或自定义指标监控批次处理速率、线程池活跃线程数、队列大小、数据库连接池使用情况等并接入PrometheusGrafana。记录失败批次的数据和原因方便后续补偿或人工干预。优雅停止与状态保存对于长时间运行的导入任务应支持优雅停止如监听应用关闭事件PreDestroy。考虑将任务进度如已处理的批次索引持久化到数据库或Redis中任务重启后可以从断点继续避免重复处理或数据丢失。数据验证与清洗在插入数据库前进行必要的数据验证非空、格式、长度、业务规则。对于脏数据如格式错误不应导致整个批次失败可以将其放入一个“错误数据”集合任务完成后统一报告或处理。连接池与线程池的精细调优连接池maximum-pool-size需要根据数据库最大连接数和应用实例数综合设定。监控数据库的Threads_connected。线程池对于I/O密集型任务如数据库操作线程数可以设置得比CPU核心数多。通过监控线程池的队列堆积情况来调整queueCapacity和maxPoolSize。生产环境事务策略慎用大事务本文多线程示例中每个批次独立事务是推荐做法。考虑最终一致性对于超大数据量可以引入消息队列如RocketMQ、Kafka将导入任务异步化实现削峰填谷和最终一致性。SQL优化对于MySQL如果导入的表是空的可以考虑先ALTER TABLE ... DISABLE KEYS禁用索引导入完成后再ALTER TABLE ... ENABLE KEYS重建速度会快很多。使用LOAD DATA INFILE命令直接从文件导入这是MySQL原生最快的方式适合运维操作或特定场景。代码健壮性使用try-with-resources确保资源如文件流、数据库连接关闭。对可能为null的集合使用CollectionUtils.isEmpty()进行判断。重要的业务方法考虑添加单元测试和集成测试模拟大数据量场景。掌握海量数据高效写入是后端开发者处理数据密集型应用的基本功。本文从问题出发逐步构建了一个包含单线程分批、多线程并行、失败重试的完整解决方案并深入探讨了连接池、线程池配置、事务边界等底层细节。关键在于理解“分治”思想并合理利用现有工具MyBatis-Plus、线程池来简化开发。在实际项目中你需要根据数据源类型文件、消息、API、数据规模、业务对一致性的要求强一致/最终一致以及运维能力灵活选择和组合这些策略。例如对于实时性要求不高的报表数据初始化可以采用多线程分批对于支付流水等强一致性要求高的场景则可能需要在事务和性能之间做出更谨慎的权衡。建议你在本地运行示例代码通过调整参数观察性能变化并尝试集成到自己的项目中。遇到具体问题时再回头查阅本文的“常见问题”和“最佳实践”部分相信能找到解决思路。
返回列表