一句话结论:future、completionservice、completablefuture 都能把串行请求变并发,但只有后两者能让你按完成顺序拿到结果。而在真实线上系统里,决定这段代码是否可靠的,往往不是选了哪个 api,而是线程池、超时预算、异常语义和可观测性这四件事有没有做对。
一、从一个真实场景说起
1.1 场景与约束
一个很常见的需求:需要从下游服务拉取一批数据,数据量大到一次拉不完,接口只提供分页。
串行分页的耗时模型是这样的:
总耗时 ≈ 页数 × 单页 rt
100 页、单页 200ms,就是 20 秒。这个耗时放在定时任务里勉强能接受,放在一个同步接口里基本不可用。
于是自然想到:既然每一页的拉取彼此独立,为什么不并发?并发之后的耗时模型变成:
总耗时 ≈ ceil(页数 / 并发度) × 单页 rt + 调度开销
1.2 真正要解决的问题:结果什么时候"可用"
批量任务的耗时几乎不可能一致:有的页命中缓存 10ms 返回,有的页穿透到 db 走了慢查询 3s 返回。这时有两个完全不同的诉求:
- 诉求 a:我要拿到全部结果再统一处理。 比如要对全量数据做排序、聚合、去重。
- 诉求 b:谁先完成我就先处理谁。 比如边拉边写库、边拉边推 mq、边拉边流式返回给前端;或者在超时预算内能拿多少算多少,剩下的降级。
诉求 a 用 future 就够了。诉求 b 用 future 会很别扭——这也是原文想说清楚的核心点。
这里要先纠正一个容易混淆的地方:"能优先拿到最快完成的结果"改善的是首个结果的可用时间和处理的流水化程度,不等于总耗时会变短。 如果代码最终还是要等所有任务结束,那么三种写法的墙钟耗时都约等于最慢任务的耗时。下面的示例会刻意把这一点暴露出来。
为了让示例可复现,先准备一个统一的辅助方法(注意 interruptedexception 的正确处理,后面会解释为什么不能只 printstacktrace):
static void sleepseconds(long seconds) {
try {
timeunit.seconds.sleep(seconds);
} catch (interruptedexception e) {
// 不要吞掉中断信号:恢复中断标记,让上层有机会感知取消
thread.currentthread().interrupt();
throw new illegalstateexception("task interrupted", e);
}
}
三个耗时不等的任务(用秒是为了在控制台肉眼可辨,生产代码里当然是毫秒级 rpc):
static list<callable<integer>> buildtasks() {
return arrays.aslist(
() -> { sleepseconds(3); system.out.println("task-3s done"); return 3; },
() -> { sleepseconds(1); system.out.println("task-1s done"); return 1; },
() -> { sleepseconds(2); system.out.println("task-2s done"); return 2; }
);
}
二、方案一:invokeall + future
2.1 代码
public static void byinvokeall() throws interruptedexception {
long start = system.currenttimemillis();
executorservice pool = executors.newfixedthreadpool(3); // 仅为示例,生产别这么写,见 6.1
try {
list<future<integer>> futures = pool.invokeall(buildtasks());
for (future<integer> f : futures) {
try {
system.out.println("result: " + f.get());
} catch (executionexception e) {
system.out.println("task failed: " + e.getcause());
}
}
} finally {
pool.shutdown();
}
system.out.println("cost=" + (system.currenttimemillis() - start) + "ms");
}
输出(顺序稳定):
task-1s done
task-2s done
task-3s done
result: 3
result: 1
result: 2
cost=3005ms
2.2 为什么拿不到"最先完成"的结果
关键不在 future,而在 invokeall 的语义:
invokeall 会阻塞调用线程,直到所有任务全部完成(或被取消),然后返回一个与入参集合迭代顺序一致的 future 列表。
也就是说,invokeall 返回的那一刻,所有 future 都已经是 isdone() == true 了。后面那个 for 循环里的 f.get() 不会再阻塞,它只是按提交顺序把已经躺在那儿的结果读出来。
所以"先完成先返回"在这个写法里 根本没有机会发生——你在拿第一个结果之前,就已经等完了最慢的那个任务。日志里 task-1s done 出现在最前面,但 result: 3 出现在最前面,这个错位正好说明了问题:任务的完成顺序和结果的消费顺序被解耦开了,而消费顺序被固定成了提交顺序。
一个常见的"自己造轮子"补救方案是轮询:
// 反面示例:忙轮询 isdone()
while (!futures.isempty()) {
futures.removeif(f -> {
if (f.isdone()) { consume(f); return true; }
return false;
});
}
这段代码能"跑通",但它会把一个 cpu 核心打满去做无意义的空转,还容易在任务数多时引入明显的调度抖动。completionservice 存在的意义就是把这件事做对。
2.3 future 的能力边界
future 是 jdk 5 的产物,它的接口只提供了四种能力:查完成、取消、阻塞取结果、带超时地阻塞取结果。它没有的能力:
- 没有完成回调(只能主动问,不能被动通知);
- 没有组合能力(无法表达"a 完成后接 b"、"a、b 都完成后合并");
- 没有按完成顺序获取的入口;
- 异常必须通过
get()抛出的executionexception去getcause()剥一层。
2.4 什么时候用它反而是对的
不要因为它"老"就否定它。以下场景 invokeall 是最简洁、最不容易写错的选择:
- 语义上就是"全部拿到才有意义",比如批量校验、全量聚合;
- 任务数量固定且不多,耗时方差不大;
- 你需要一个整体超时并自动取消未完成任务——这时用带超时的重载,非常好用:
// 3 秒整体预算,到点未完成的任务会被 cancel
list<future<integer>> futures = pool.invokeall(buildtasks(), 3, timeunit.seconds);
for (future<integer> f : futures) {
if (f.iscancelled()) {
system.out.println("timeout, skipped"); // 超时被取消
continue;
}
try {
system.out.println("result: " + f.get());
} catch (executionexception e) {
system.out.println("failed: " + e.getcause());
}
}
这个重载常被忽略,但它把"整体超时 + 未完成任务取消"两件事一次做了,比自己拿 countdownlatch 拼要稳妥。
需要留意的是:cancel(true) 只是给线程发中断信号。如果任务体内部是不响应中断的阻塞(例如某些 http 客户端只认自己的 socket timeout),线程不会立刻退出。取消是"请求取消",不是"保证停止",这一点在排查"任务超时了但线程池还满着"时非常关键。
三、方案二:completionservice
3.1 代码
public static void bycompletionservice() {
long start = system.currenttimemillis();
executorservice pool = executors.newfixedthreadpool(3);
try {
list<callable<integer>> tasks = buildtasks();
// 注意:completionservice 要每批次新建,不要做成共享单例,见 3.4
completionservice<integer> cs = new executorcompletionservice<>(pool);
tasks.foreach(cs::submit);
for (int i = 0; i < tasks.size(); i++) {
try {
integer r = cs.take().get(); // 按完成顺序返回
system.out.println("result: " + r);
} catch (executionexception e) {
system.out.println("task failed: " + e.getcause());
} catch (interruptedexception e) {
thread.currentthread().interrupt();
break;
}
}
} finally {
pool.shutdown();
}
system.out.println("cost=" + (system.currenttimemillis() - start) + "ms");
}
输出:
task-1s done
result: 1
task-2s done
result: 2
task-3s done
result: 3
cost=3006ms
结果按完成顺序出现,第一个结果在 1 秒时就已经可以处理了。总耗时依然约 3 秒——因为循环还是要取满三个结果。这正好印证 1.2 节的判断。
3.2 原理:多了一个完成队列
executorcompletionservice 的实现思路很直白,没有任何魔法:它把你提交的 callable 包装成一个 queueingfuture,重写了 futuretask 的 done() 钩子——任务完成时,把自己塞进一个内部的 blockingqueue。
于是:
submit():提交到底层executor,同时登记到完成队列;take():从完成队列阻塞取一个已完成的future;poll()/poll(timeout, unit):非阻塞 / 限时获取。
理解了这一点,就能理解它的所有约束:它的能力完全来自那个队列,它不知道你提交了多少个任务。
3.3 关键用法:整体超时预算 + 部分成功
completionservice 真正的工程价值不在"顺序好看",而在于它天然适合"在给定的时间预算内,能拿多少算多少"。这是接口层做并发聚合时最实用的一个模式:
public class batchfanout {
public static class batchresult<t> {
public final list<t> succeeded = new arraylist<>();
public int failed;
public int unfinished;
public boolean ispartial() { return failed > 0 || unfinished > 0; }
}
/**
* 在 budgetms 的整体预算内并发执行 tasks,返回已完成的结果,未完成的取消。
*/
public static <t> batchresult<t> execute(list<callable<t>> tasks,
executorservice pool,
long budgetms) {
batchresult<t> result = new batchresult<>();
completionservice<t> cs = new executorcompletionservice<>(pool);
list<future<t>> submitted = new arraylist<>(tasks.size());
for (callable<t> task : tasks) {
submitted.add(cs.submit(task));
}
long deadline = system.nanotime() + timeunit.milliseconds.tonanos(budgetms);
try {
for (int i = 0; i < tasks.size(); i++) {
long remaining = deadline - system.nanotime();
future<t> f = remaining <= 0 ? null : cs.poll(remaining, timeunit.nanoseconds);
if (f == null) { // 预算用尽
result.unfinished = tasks.size() - i;
break;
}
try {
result.succeeded.add(f.get()); // 此处不会再阻塞
} catch (executionexception e) {
result.failed++;
// 生产代码请打日志并带上任务标识,见第七章
}
}
} catch (interruptedexception e) {
thread.currentthread().interrupt();
} finally {
// 必须收尾:提前 break 后,未完成任务仍在占用线程池
submitted.foreach(f -> f.cancel(true));
}
return result;
}
}
这段代码里有三个容易漏掉的细节,值得单独强调:
- 超时预算是"整体 deadline",不是"每个任务 timeout"。 如果给每个任务各自 500ms,10 个任务串在一起最坏就是 5 秒;用整体 deadline,总耗时才真正可控。上游接口的 sla 是整体的,超时预算也应该是整体的。
- 提前退出后一定要
cancel剩余任务。 否则这批任务会继续占用线程池,下一次请求进来时线程池已经被上一批的"僵尸任务"占满——这是并发聚合类代码最常见的雪崩诱因。 f.get()放在poll之后是安全的。 从完成队列里取出来的future一定已经 done,get()不会阻塞。
3.4 有哪些坑
| 坑 | 现象 | 正确做法 |
|---|---|---|
take() 次数多于实际提交数 | 线程永久阻塞,接口一直不返回 | 严格用提交数计数;能提前退出的场景用 poll(timeout) |
把 completionservice 做成共享单例/成员变量 | 请求 a 取到了请求 b 的结果,数据串批 | 每批次新建一个实例,只共享底层 executorservice |
提前 break 却不 cancel | 线程池被上一批任务占满,后续请求排队甚至拒绝 | finally 中统一 cancel(true) |
submit 抛 rejectedexecutionexception 后继续 take | 提交数与计数不一致 → 阻塞 | 提交在 try 内,用实际成功提交的数量作为循环上界 |
吞掉 interruptedexception | 上层取消/超时后线程仍在跑 | thread.currentthread().interrupt() 后退出 |
completionservice 的能力边界也很清楚:它解决了"按完成顺序消费",但不解决编排。"a 完成后触发 b,b、c 完成后合并成 d"这类依赖关系,用它写出来就是一堆嵌套的 take 循环。这就是 completablefuture 的地盘。
四、方案三:completablefuture
4.1 代码
completablefuture 提供的是回调式模型:结果就绪时由完成它的线程回调你的处理逻辑,天然就是"先完成先处理"。
public static void bycompletablefuture(executorservice pool) {
long start = system.currenttimemillis();
list<integer> seconds = arrays.aslist(3, 1, 2);
list<completablefuture<integer>> futures = seconds.stream()
.map(s -> completablefuture.supplyasync(() -> {
sleepseconds(s);
system.out.println("task-" + s + "s done on " + thread.currentthread().getname());
return s;
}, pool)) // 显式指定线程池,见 4.2
.collect(collectors.tolist());
// 完成即消费(谁先完成谁先进来)
futures.foreach(f -> f.thenaccept(r -> system.out.println("result: " + r)));
// 等全部结束;任一异常会在这里抛 completionexception
completablefuture.allof(futures.toarray(new completablefuture[0])).join();
system.out.println("cost=" + (system.currenttimemillis() - start) + "ms");
}
输出:
task-1s done on batch-pool-2
result: 1
task-2s done on batch-pool-3
result: 2
task-3s done on batch-pool-1
result: 3
cost=3011ms
4.2 第一个坑:不要用默认线程池
原文示例用的是 supplyasync(supplier) 单参版本。它的执行器是 forkjoinpool.commonpool(),这在业务代码里通常是错的,原因有三个:
其一,它是整个 jvm 共享的。 你的批量拉取、别人的 parallelstream()、某个库内部的异步逻辑,全都挤在同一个池子里。一段慢的阻塞任务会把不相关的功能一起拖死,而且极难定位——线程栈上只会看到 forkjoinpool.commonpool-worker-n,看不出是谁提交的。
其二,它的默认并行度是 可用处理器数 - 1。 在 4 核容器里就是 3。这个池子被设计用来跑 cpu 密集的计算任务,而不是跑一堆 io 阻塞任务。更隐蔽的是:当计算出的并行度为 0 时(例如容器只分到 1 核),commonpool 会退化成"每个任务新建一个线程"的模式。这意味着你的并发行为会随部署环境的 cpu 配额而变——同一份代码在 8 核机器和 1 核容器里表现完全不同。
其三,它不响应你的治理手段。 你没法给它设置有界队列、拒绝策略、监控指标和有意义的线程名。
正确做法是自己建池并显式传入:
threadpoolexecutor batchpool = new threadpoolexecutor(
8, 8,
60l, timeunit.seconds,
new arrayblockingqueue<>(200), // 有界
new customizablethreadfactory("batch-fetch-"), // spring 提供,线程名可辨识
new threadpoolexecutor.abortpolicy()); // 快速失败优于悄悄堆积
batchpool.allowcorethreadtimeout(true);
顺便说一句:parallelstream() 同样跑在 commonpool 上,且不支持指定执行器。在业务代码里对 io 操作用 parallelstream(),属于同一类问题。
4.3 第二个坑:join 和 get 的异常语义不一样
这是 code review 里高频出现的混淆点:
| 方法 | 抛出的异常 | 是否受检 |
|---|---|---|
get() | executionexception、interruptedexception | 是,必须处理 |
join() | completionexception | 否,运行时异常 |
两者包装的原始异常都要通过 getcause() 取。join() 写起来干净(尤其在 lambda 里),代价是编译器不会提醒你处理失败。原文示例中 futures.foreach(completablefuture::join) 如果第一个任务抛异常,整个方法直接抛出,后面的任务既没被消费也没被取消——异常路径上的行为和正常路径完全不同。
另外,thenaccept 注册的回调不会捕获上游异常。上游失败时,thenaccept 会被跳过,异常沿链向下传递。要处理失败必须显式接上 exceptionally / handle / whencomplete:
completablefuture.supplyasync(() -> client.fetchpage(pageno), batchpool)
.thenaccept(this::consume) // 成功路径
.exceptionally(ex -> { // 失败路径,返回兜底值让链继续
log.warn("fetch page {} failed", pageno, ex);
return null;
});
三者的区别:exceptionally 只在异常时触发并可替换结果;handle 无论成功失败都触发且可改结果;whencomplete 无论成功失败都触发但不改变结果(异常会继续传递)。
4.4 组合能力:allof / anyof / 超时
这是 completablefuture 相对前两者真正的增量价值。
全部完成后合并(注意 allof 返回 completablefuture<void>,结果要自己捞):
completablefuture<list<pageresult>> all = completablefuture
.allof(futures.toarray(new completablefuture[0]))
.thenapply(v -> futures.stream()
.map(completablefuture::join) // 此时都已完成,join 不阻塞
.collect(collectors.tolist()));
任一完成即返回(典型场景:多个数据源竞速取最快的那个):
completablefuture<object> fastest = completablefuture.anyof(fromcache, fromdb);
超时(java 9+):
completablefuture<pageresult> f = completablefuture
.supplyasync(() -> client.fetchpage(pageno), batchpool)
.ortimeout(500, timeunit.milliseconds) // 超时则异常完成
.exceptionally(ex -> pageresult.empty(pageno)); // 降级为空页
// 或者直接给默认值,不走异常路径
completablefuture<pageresult> f2 = completablefuture
.supplyasync(() -> client.fetchpage(pageno), batchpool)
.completeontimeout(pageresult.empty(pageno), 500, timeunit.milliseconds);
这里有一个必须知道的语义细节:ortimeout / completeontimeout 只是让这个 completablefuture 提前进入完成状态,它不会中断正在执行的那个任务。任务体会继续跑到自然结束,继续占着线程池的一个位置。也就是说,超时保护的是调用方的响应时间,不保护后端资源。真正的资源保护要靠:任务内部自己的 socket/read timeout、连接池上限、以及有界队列 + 拒绝策略。
如果你用 java 8,没有 ortimeout,就需要自己拿 scheduledexecutorservice 去 completeexceptionally,或者干脆用 completionservice 的 poll(timeout) 模式(3.3 节),后者在 java 8 项目里通常更简单。
4.5 第三个坑:回调在哪个线程执行
非 async 后缀的方法(thenaccept、thenapply…)的执行线程是不确定的:
- 如果注册时上游尚未完成,回调由完成上游的那个线程执行;
- 如果注册时上游已经完成,回调由当前调用线程执行。
这带来两个实际后果:
- 在回调里做重活或阻塞操作(写库、发 mq、再发一次 rpc),会占用业务线程池的工作线程,甚至可能占用主线程/web 容器线程。要隔离就用
thenacceptasync(action, anotherpool),把不同性质的工作放到不同的池子里。 - 调试时看到的线程名可能和你预期完全不同,别据此推断"任务跑在哪个池"。
一个容易踩的死锁场景:在同一个固定大小的线程池里,任务 a 内部又提交任务 b 并 join() 等待 b。当池被 a 类任务占满时,b 永远排不上队,a 永远等不到 b。不要在同一个池内做嵌套等待——这是"接口偶发性完全无响应、线程栈全是 waiting"这类问题的典型成因。
五、三种方案怎么选
5.1 对比表
| 维度 | invokeall + future | completionservice | completablefuture |
|---|---|---|---|
| 结果获取顺序 | 提交顺序 | 完成顺序 | 完成顺序(回调) |
| 编程模型 | 阻塞式 | 阻塞式(队列驱动) | 回调式 / 声明式 |
| 首个结果可用时间 | 等最慢任务 | 最快任务完成即可用 | 最快任务完成即可用 |
| 任务编排(串/并/合并) | 不支持 | 不支持 | 支持(thencompose / thencombine / allof / anyof) |
| 整体超时 | invokeall(tasks, timeout, unit) | poll(timeout) + 手动 cancel | ortimeout / completeontimeout(java 9+) |
| 异常处理 | executionexception.getcause() | 同左,可按任务粒度捕获 | 链式:exceptionally / handle / whencomplete |
| 默认线程池 | 必须自己传 | 必须自己传 | 单参重载用 commonpool(陷阱) |
| 心智负担 | 低 | 中 | 中偏高(回调线程、异常传播) |
| 最低版本 | java 5 | java 5 | java 8(超时 api 需 9+) |
5.2 选型判断
按"业务语义"而不是"api 新旧"来选:
- 只要全量结果、任务不多、耗时接近 →
invokeall。代码最短,最不容易写错。需要整体超时就用带 timeout 的重载。 - 需要边完成边处理,或需要"预算内能拿多少算多少"的部分成功 →
completionservice。它是这个场景下最直白的工具,尤其适合 java 8 项目。 - 任务之间有依赖、需要编排,或要接入响应式/异步链路 →
completablefuture。前提是你能管住线程池和异常传播。 - spring 项目里的常规异步 →
@async配合自定义threadpooltaskexecutor,返回completablefuture。注意@async走的是代理,同类内部方法直接调用不会生效,这是最常被问的一个坑。
反过来,有几种情况不适合做批量并发:
- 下游是共享资源且没有做隔离(同一个 db 实例、同一个 redis 分片)。你把串行改并发,只是把压力从自己身上转移到了下游,很可能换来一批慢查询和连接池耗尽。
- 任务之间有强顺序或事务语义。并发会破坏顺序性,且跨线程后本地事务不共享——
@transactional不会传播到异步线程,回滚边界会碎掉。 - 单个任务本身很快(微秒级),任务数极多。线程调度和上下文切换开销可能超过收益,批量合并请求(batch api)通常比并发更有效。
六、工程化落地:让这段代码上线后不出事
前面讲的是 api,这一章讲的是"为什么线上出问题"。
6.1 不要用 executors 的工厂方法
executors.newfixedthreadpool(n) 用的是无界 linkedblockingqueue。下游变慢时,任务不会被拒绝,而是无声地在队列里堆积,直到内存被占满触发频繁 gc 乃至 oom。newcachedthreadpool 是另一个极端:线程数上限接近 integer.max_value,下游变慢时会疯狂创建线程,直到无法再创建原生线程。
结论:用 threadpoolexecutor 构造函数显式声明七个参数,队列必须有界。这也是《阿里巴巴 java 开发手册》里明确要求的一条。
线程数怎么定?给一个可用的起点,而不是精确公式:
- cpu 密集型:接近核心数;
- io 密集型:
核心数 × (1 + 平均等待时间 / 平均计算时间)是一个粗略上界,实践中更常见的做法是先按"期望 qps × 平均 rt"估算并发量,再压测调整; - 无论哪种,最终要靠压测和线上指标校准,公式只是初值。
拒绝策略的选择也有讲究:callerrunspolicy 常被当作"温柔的降级",但它会让提交者线程去执行任务。如果提交者是 tomcat 的工作线程,你就是在用 http 线程去跑批量任务——反压确实产生了,但反压对象是你的接口吞吐量。对于可丢弃的旁路任务,abortpolicy + 记录指标 + 告警通常更可控。
6.2 超时预算要"自上而下"分配
一个接口的总 rt 预算(比如 1 秒)应该向下拆解:
接口预算 1000ms
├── 参数校验 + 本地缓存 ~10ms
├── 批量并发拉取 预算 600ms ← 整体 deadline,不是单任务 timeout
└── 聚合、序列化、返回 剩余
批量拉取拿到的是 600ms 的整体预算。落到代码上就是 3.3 节的 deadline 模式。同时别忘了任务内部也要有超时——http 客户端的 connecttimeout / sockettimeout、数据库的 querytimeout。上层超时管的是响应时间,下层超时管的是资源释放,两者不能互相替代。
6.3 明确"部分成功"是什么语义
批量并发一旦引入超时和降级,返回结果就可能是不完整的。这时必须由业务来回答三个问题,而不是让代码默默决定:
- 部分成功算成功还是失败?返回体里要不要带
partial标记和缺失范围? - 缺失的数据是补拉(异步重试 / 落到延时队列)还是直接丢?
- 调用方能否感知?如果调用方拿到一个"看起来完整"的残缺列表,后续的聚合、对账、库存计算都可能出错。
这类问题不解决,"并发优化"就会变成"偶发性数据不一致"。在我的经验里,批量并发上线后的事故,绝大多数不是并发本身写错了,而是部分成功的语义没定义清楚。
6.4 保护下游
- 限制并发度:并发度不是越大越好,要匹配下游的承载能力。除了控制线程池大小,也可以用
semaphore在任务内部做二次限流。 - 给下游连接池留余量:8 个线程并发查库,就意味着瞬时占用 8 个数据库连接。线程池大小超过连接池大小时,多出来的线程只是在等连接。
- 失败要有边界:熔断/隔离(sentinel、resilience4j 等)应该作用在下游调用上,避免一个慢下游把整个批量任务的线程池长期占满。
七、可观测性与线上排查
并发代码最大的成本不在写,在出问题时看不见。
7.1 线程名是排查的起点
默认线程名 pool-1-thread-3 没有任何信息量。给每个池起一个业务可辨识的名字,抓一次 thread dump 就能立刻定位是谁在阻塞:
// spring 项目
new customizablethreadfactory("order-batch-fetch-");
// 或用 jdk 原生方式
threadfactory factory = r -> {
thread t = new thread(r, "order-batch-fetch-" + counter.incrementandget());
t.setdaemon(false);
t.setuncaughtexceptionhandler((th, e) -> log.error("uncaught in {}", th.getname(), e));
return t;
};
顺带说一句 submit 和 execute 的差别:用 submit 提交时,任务抛出的异常会被封进 future,如果你不调用 get(),异常就被彻底吞掉了,uncaughtexceptionhandler 也不会触发。这是"任务好像没执行,但也没有任何报错日志"的常见原因。
7.2 traceid 必须跨线程透传
异步化之后,子线程的日志会丢掉 mdc 里的 traceid,链路直接断在这里。spring 的 threadpooltaskexecutor 提供了 taskdecorator 扩展点:
public class mdctaskdecorator implements taskdecorator {
@override
public runnable decorate(runnable runnable) {
map<string, string> parent = mdc.getcopyofcontextmap(); // 提交线程的上下文
return () -> {
try {
if (parent != null) {
mdc.setcontextmap(parent);
}
runnable.run();
} finally {
mdc.clear(); // 池化线程必须清理,防止污染下一个任务
}
};
}
}
@bean("batchexecutor")
public threadpooltaskexecutor batchexecutor() {
threadpooltaskexecutor executor = new threadpooltaskexecutor();
executor.setcorepoolsize(8);
executor.setmaxpoolsize(8);
executor.setqueuecapacity(200);
executor.setthreadnameprefix("batch-fetch-");
executor.settaskdecorator(new mdctaskdecorator());
executor.initialize();
return executor;
}
注意 finally 里的清理:线程是复用的,不清理会导致下一个任务打出上一个请求的 traceid,比没有 traceid 更容易误导排查。
如果需要透传的不只是 mdc,还有自定义的 threadlocal 上下文(用户身份、灰度标记、租户 id),可以考虑 transmittablethreadlocal(alibaba/transmittable-thread-local),它提供了 ttlexecutors 来包装线程池。
7.3 线程池要有指标
线程池是黑盒,必须暴露出来。threadpoolexecutor 自带这些可读属性:
| 方法 | 含义 | 关注点 |
|---|---|---|
getactivecount() | 正在执行任务的线程数 | 长期贴着 maxpoolsize → 容量不足或任务变慢 |
getqueue().size() | 队列积压 | 持续增长 → 下游变慢,即将触发拒绝 |
getcompletedtaskcount() | 累计完成数 | 增速停滞 → 任务卡死 |
getlargestpoolsize() | 历史峰值线程数 | 判断 maxpoolsize 是否被打满过 |
接入 micrometer 只需要一行包装,之后线程池指标就能进 prometheus / grafana:
executorservice monitored = executorservicemetrics.monitor(
meterregistry, batchpool, "batch-fetch", tags.of("biz", "order"));
再补三个业务维度的指标,排查效率会完全不同:批次总耗时、单任务耗时分布(p99)、部分成功率 / 超时任务数。
7.4 排查对照表
| 线上现象 | 优先看什么 | 常见根因 |
|---|---|---|
| 接口 rt 从毫秒级跳到超时 | 线程池队列长度、activecount | 队列积压,下游变慢;线程池被占满 |
| 接口偶发完全无响应,线程栈大量 waiting | thread dump 中的 park 位置 | 同池嵌套等待死锁(4.5);take() 次数多于提交数(3.4) |
| 任务像没执行,但没有异常日志 | 是否用 submit 且从未 get() | 异常被封进 future 后吞掉(7.1) |
| 老年代持续增长直至 oom | 队列类型是否无界 | executors.newfixedthreadpool 的无界队列(6.1) |
| 日志里 traceid 丢失或串了 | 是否配置 taskdecorator 且 finally 清理 | mdc 未透传 / 未清理(7.2) |
| 超时已生效但线程池仍满 | 任务是否响应中断、下游 socket timeout | ortimeout 不中断任务(4.4);cancel 只是发中断(2.4) |
| 同一份代码在不同环境并发表现不一致 | 是否用了 commonpool / parallelstream | 并行度随 cpu 配额变化(4.2) |
八、总结
回到原文的那个问题——"如何让先完成的任务先返回结果",答案分三层:
- api 层面:
invokeall的语义就是"等全部完成",所以它做不到;completionservice通过完成队列做到了;completablefuture通过回调做到了,并且额外提供了编排能力。 - 收益层面:按完成顺序消费改善的是首个结果的可用时间和处理的流水化程度,不等于总耗时变短。只有配合超时预算、提前退出、流式输出,它才会转化为可感知的 rt 收益。
- 可靠性层面:这段代码上线后能不能扛住,取决于线程池是否有界、超时预算是否自上而下、部分成功语义是否定义清楚、traceid 和线程池指标是否可观测。这四件事比选哪个 api 重要得多。
到此这篇关于java批量并发请求的三种写法与性能优化的文章就介绍到这了,更多相关java批量并发请求内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论