如何用ScheduledExecutorService替代Thread.sleep分批处理大列表并终止任务
使用ScheduledExecutorService实现分批延迟处理并优雅终止
第一步:拆分大列表为固定大小的批次
先实现一个工具方法,把输入的大列表拆分成每10个元素一组的子批次:
private static <T> List<List<T>> splitIntoBatches(List<T> inputList, int batchSize) { List<List<T>> batches = new ArrayList<>(); for (int i = 0; i < inputList.size(); i += batchSize) { int endIndex = Math.min(i + batchSize, inputList.size()); batches.add(inputList.subList(i, endIndex)); } return batches; }
第二步:用ScheduledExecutorService实现带延迟的批次调度
这里采用递推调度的逻辑:处理完当前批次后,再调度下一批次在指定延迟后执行,完美替代Thread.sleep的阻塞式延迟,同时保证线程安全和优雅终止:
public void processLargeList(List<YourData> largeList) { int batchSize = 10; long delayMs = 60_000; // 批次间隔60秒 List<List<YourData>> batches = splitIntoBatches(largeList, batchSize); if (batches.isEmpty()) return; // 创建单线程调度器,保证批次按顺序执行 ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); AtomicInteger currentBatchIdx = new AtomicInteger(0); Runnable batchTask = new Runnable() { @Override public void run() { int idx = currentBatchIdx.getAndIncrement(); // 所有批次处理完成,关闭调度器 if (idx >= batches.size()) { executor.shutdown(); try { // 等待1分钟确保剩余任务完成,超时则强制终止 if (!executor.awaitTermination(1, TimeUnit.MINUTES)) { executor.shutdownNow(); } } catch (InterruptedException e) { executor.shutdownNow(); Thread.currentThread().interrupt(); } return; } // 处理当前批次 List<YourData> batch = batches.get(idx); try { yourService.processBatch(batch); // 替换为你的实际服务方法 } catch (Exception e) { // 根据业务需求处理异常:日志记录、跳过或终止后续批次 System.err.printf("批次 %d 处理失败: %s%n", idx, e.getMessage()); // 若要终止后续批次,直接return即可 // return; } // 调度下一个批次(还有剩余的话) if (currentBatchIdx.get() < batches.size()) { executor.schedule(this, delayMs, TimeUnit.MILLISECONDS); } } }; // 立即执行第一个批次 executor.submit(batchTask); }
关键细节说明
- 调度器选择:用
newSingleThreadScheduledExecutor()保证批次按顺序执行,避免并发处理的冲突;如果业务允许并发处理批次,可改用多线程调度器,但需自行处理线程安全问题。 - 线程安全计数:用
AtomicInteger跟踪当前批次索引,避免多线程环境下的计数错误。 - 优雅终止逻辑:
- 所有批次处理完成时,调用
executor.shutdown()拒绝新任务,等待已提交任务执行完毕。 awaitTermination()用于等待调度器关闭,超时则调用shutdownNow()强制终止,防止资源泄漏。
- 所有批次处理完成时,调用
- 与Thread.sleep的差异:这种异步调度方式不会阻塞API请求线程,批次处理在后台执行,API可快速返回响应;同时任务的调度、取消逻辑更可控,不会出现
Thread.sleep难以中断的问题。
可选:提前调度所有批次的方式
如果你的批次处理时间稳定,且允许批次按固定时间点执行(不管前一批是否完成),可以提前计算每个批次的延迟时间,一次性提交所有任务:
public void scheduleAllBatches(List<YourData> largeList) { int batchSize = 10; long delayMs = 60_000; List<List<YourData>> batches = splitIntoBatches(largeList, batchSize); if (batches.isEmpty()) return; ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); for (int i = 0; i < batches.size(); i++) { List<YourData> batch = batches.get(i); // 第i个批次延迟 i * delayMs 执行 executor.schedule(() -> { try { yourService.processBatch(batch); } catch (Exception e) { System.err.printf("批次 %d 处理失败: %s%n", i, e.getMessage()); } // 最后一个批次执行完后关闭调度器 if (i == batches.size() - 1) { executor.shutdown(); // 可添加awaitTermination逻辑保证关闭完成 } }, i * delayMs, TimeUnit.MILLISECONDS); } }
注意:这种方式若前一批处理时间超过60秒,下一批会在预定时间执行,可能导致批次重叠,仅适合处理时间稳定的场景。
内容的提问来源于stack exchange,提问作者Muhammad Umer
相关产品推荐
相关产品推荐

