linkedblockingqueue 配合 completablefuture 实现异步批处理,核心在于用队列缓冲任务、用 completablefuture 控制异步执行与结果聚合,避免阻塞主线程又保证批量处理的可控性。

用 LinkedBlockingQueue 缓存待处理任务
LinkedBlockingQueue 是线程安全的有界/无界阻塞队列,适合在多线程环境下暂存待批处理的任务(如请求对象、数据记录等)。它天然支持生产者-消费者模型,可让上游快速入队、下游按需拉取批次。
- 定义队列时建议指定容量(如
new LinkedBlockingQueue(1000)),防止内存无限增长 - 上游调用
queue.offer(task)非阻塞入队;若需等待空间可用,可用put()(但注意可能阻塞) - 避免直接在 offer 后立即触发处理——要由独立的“批处理器”线程统一拉取
启动后台批处理线程,定期拉取并提交 CompletableFuture
单独起一个守护线程(或复用线程池),定时/按量从队列中拉取任务,封装为 CompletableFuture 提交到业务线程池执行。
Java JDK 25 来自 OpenJDK 官方归档,版本为 JDK 25,本条下载地址已指向官方 Windows x64 zip 安装包直链,适合调试旧项目或兼容旧版 Java 运行环境。
- 推荐用
ScheduledExecutorService每 100ms 检查一次:若队列有 ≥50 个任务,就批量拉取(drainTo(list, batchSize)) - 对每批任务创建一个
CompletableFuture.supplyAsync(() -> processBatch(list), executor) - 注意:不要在 supplyAsync 内部再用
get()或join()阻塞,保持异步链路畅通
用 allOf 聚合批次结果,统一回调或异常处理
当一批任务被拆成多个 CompletableFuture 并行执行时,可用 CompletableFuture.allOf() 等待全部完成,再用 thenApply() 收集结果。
- 示例:收集一批用户 ID 的详细信息,每个 ID 对应一个
CompletableFuture<user></user>,最后用allOf(futures).thenApply(v -> Stream.of(futures).map(Future::join).collect(...)) - 更健壮的做法是用
thenCompose+handle处理单个 future 的失败,避免一个失败导致整批中断 - 结果可发回监听器、写入数据库、或通过回调函数通知原始调用方
配合 CompletionStage 实现任务链式编排
如果批处理后还需后续操作(如写日志、触发通知、更新状态),可把整个流程串成 CompletionStage 链,提升可读性和错误隔离能力。
- 例如:
batchFuture.thenApply(this::enrichResult).thenAcceptAsync(this::persistToDB, ioExecutor) - 每个环节都返回新的 CompletableFuture,异常可通过
exceptionally()或handle()捕获并降级处理 - 避免在链中混用同步逻辑(如耗时 IO),否则会拖慢整个异步流
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南










