Java分页+多线程处理大数据集异常及正确实现方案
海量数据库数据分页+多线程异步处理的正确实现方式
核心实现思路
- 分页读取数据库数据,避免一次性加载全量数据引发内存溢出(OOM)
- 用自定义线程池管理异步任务,严格控制并发数,防止压垮数据库或系统资源
- 统一管理任务生命周期,确保所有异步任务执行完成后再结束流程
- 全链路异常捕获,避免单个任务失败导致整个流程中断
完整代码示例
1. 线程池初始化
import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; // 自定义线程池,根据系统资源和数据库承载能力调整参数 private static final ExecutorService dataProcessExecutor = new ThreadPoolExecutor( 4, // 核心线程数:IO密集型任务建议设为CPU核心数的2-4倍 8, // 最大线程数:不超过数据库连接池最大连接数 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(100), // 有界任务队列,避免无界队列引发OOM new ThreadFactory() { private final AtomicInteger threadCounter = new AtomicInteger(1); @Override public Thread newThread(Runnable r) { return new Thread(r, "data-handler-thread-" + threadCounter.getAndIncrement()); } }, new ThreadPoolExecutor.CallerRunsPolicy() // 队列满时由提交线程执行,避免任务丢失 );
2. 分页+多线程处理主逻辑
import java.util.ArrayList; import java.util.List; import java.util.concurrent.Future; import org.springframework.data.domain.Page; import org.springframework.data.domain.PageRequest; import org.springframework.data.domain.PageImpl; import org.slf4j.Logger; import org.slf4j.LoggerFactory; private static final Logger log = LoggerFactory.getLogger(YourClass.class); public void batchProcessLargeData() { int currentPage = 1; final int PAGE_SIZE = 1000; // 分页大小:根据单条数据大小和内存调整 List<Future<?>> taskFutures = new ArrayList<>(); try { Page<?> dataPage; do { // 分页查询数据库数据(替换为实际业务的分页查询逻辑) dataPage = queryDataByPage(currentPage, PAGE_SIZE); if (dataPage.getContent().isEmpty()) { break; } // 提交当前页数据的处理任务到线程池 List<?> currentPageData = dataPage.getContent(); taskFutures.add(dataProcessExecutor.submit(() -> { try { // 执行分析逻辑 analyzeData(currentPageData); // 执行处理逻辑 processData(currentPageData); } catch (Exception e) { // 捕获任务内所有异常,避免异常扩散 log.error("分页数据处理失败,页码:{}", currentPage, e); } })); currentPage++; } while (dataPage.hasNext()); // 等待所有异步任务执行完成 for (Future<?> future : taskFutures) { try { future.get(); } catch (InterruptedException | ExecutionException e) { log.error("等待任务执行完成异常", e); // 可根据业务需求决定是否中断整个流程 } } } finally { // 优雅关闭线程池 dataProcessExecutor.shutdown(); try { // 等待线程池完成所有任务,超时则强制关闭 if (!dataProcessExecutor.awaitTermination(1, TimeUnit.HOURS)) { dataProcessExecutor.shutdownNow(); } } catch (InterruptedException e) { dataProcessExecutor.shutdownNow(); Thread.currentThread().interrupt(); } } } // 分页查询数据库的业务方法(示例) private Page<?> queryDataByPage(int pageNum, int pageSize) { // 替换为实际实现:比如JPA的Pageable查询、MyBatis分页插件查询 return new PageImpl<>(List.of(), PageRequest.of(pageNum, pageSize), 0); } // 数据分析方法(示例) private void analyzeData(List<?> dataList) { // 实际业务的数据分析逻辑 } // 数据处理方法(示例) private void processData(List<?> dataList) { // 实际业务的数据处理逻辑 }
常见问题排查与优化
- 线程池参数调优:核心线程数根据任务类型调整(CPU密集型设为CPU核心数+1,IO密集型设为2-4倍核心数);队列必须用有界队列,避免无界队列引发内存溢出;拒绝策略根据业务选择,
CallerRunsPolicy适合不允许丢失任务的场景。 - 数据库连接池适配:确保数据库连接池的最大连接数≥线程池最大线程数,避免分页查询时连接耗尽抛出异常。
- 线程安全保障:如果
analyzeData或processData中使用了共享变量/资源,必须加锁或使用线程安全的工具类,避免并发问题。 - 分页逻辑正确性:建议用唯一主键(如ID)进行分页(代替offset分页),避免因数据新增/删除导致的重复读取或遗漏数据。
- 异常全链路处理:任务内部必须捕获所有异常,否则
Future.get()会抛出ExecutionException,需在等待任务时统一处理。
内容的提问来源于stack exchange,提问作者sriram_aleti
相关产品推荐
相关产品推荐

