十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

CompletableFuture 源码

CompletableFuture 源码 CompletableFuturethenApplypublicUCompletableFutureUthenApply(Function?superT,?extendsUfn){returnuniApplyStage(null,fn);}uniApplyStage创建新的 CompletableFuture传入的线程池不为 null 或者同步执行回调成功则直接返回否则将任务压入上游的链表中并执行回调任务privateVCompletableFutureVuniApplyStage(Executore,Function?superT,?extendsVf){if(fnull)thrownewNullPointerException();// // 创建一个新的 CompletableFuture作为本次链式调用的返回值CompletableFutureVdnewCompletableFutureV();// e ! null传入了自定义线程池 → 必须走异步任务// !d.uniApply(this, f, null)e 是 null尝试同步直接完成计算if(e!null||!d.uniApply(this,f,null)){UniApplyT,VcnewUniApplyT,V(e,d,this,f);// 将任务 c 压入上游 Future (this) 的无锁回调链表push(c);// 执行这个回调任务SYNC 代表同步模式c.tryFire(SYNC);}returnd;}uniApply判断上游是否已经执行完成如果没有执行完直接返回 false判断上游是否抛出异常如果抛出异常则将异常传入给下游任务通过 CAS 修改任务状态修改成功则执行下游任务并将结果封装为 Result// a - 上游任务// f - 下游回调函数finalSbooleanuniApply(CompletableFutureSa,Function?superS,?extendsTf,UniApplyS,Tc){Objectr;Throwablex;// a null上游 future 为空// (r a.result) null上游还没有完成// f null转换函数为 nullif(anull||(ra.result)null||fnull)returnfalse;// 当前任务未完成有可能别的线程已经把下游完成了那就不要再重复执行tryComplete:if(resultnull){if(rinstanceofAltResult){if((x((AltResult)r).ex)!null){// 上游执行抛出异常把异常封装为 Result 传给下游 futurecompleteThrowable(x,r);// 跳出标签 if不再执行后面的 Function函数不会被执行异常直接向下传递breaktryComplete;}rnull;}try{// c ! null说明当前是从UniApply.tryFire()调用过来的传入了任务对象// claim 返回 false抢占失败别的线程已经在跑这个任务if(c!null!c.claim())returnfalse;Ss(S)r;// 执行用户传入的转换函数并 CAS 把函数执行结果设置到当前 thiscompleteValue(f.apply(s));}catch(Throwableex){completeThrowable(ex);}}returntrue;}push如果当前 future 没有完成并且将任务压入栈finalvoidpush(UniCompletion?,?c){if(c!null){// result null 当前 future 还没有完成// tryPushStack(c) 执行 CAS 尝试把 c 压入回调栈CAS 失败返回 false成功返回 truewhile(resultnull!tryPushStack(c))// 把当前待压入任务 c 的 next 置为 nulllazySetNext(c,null);}}tryPushStack将 stack 指向 c并将 c 的 next 指向之前的任务finalbooleantryPushStack(Completionc){Completionhstack;lazySetNext(c,h);returnUNSAFE.compareAndSwapObject(this,STACK,h,c);}postFire同步模式则清理上游 future 的 stack 中已完成的任务异步模式则调用 postComplete 运行任务finalCompletableFutureTpostFire(CompletableFuture?a,intmode){// a ! null a.stack ! null上游 a 还有未清理完的回调链表stackif(a!nulla.stack!null){// mode 0SYNC同步模式// a.result null上游 a 还没有完成if(mode0||a.resultnull)// 清理上游 a 的 stack 链表移除执行完成的任务a.cleanStack();else// 异步模式并且 a 已经执行完成// postComplete 消费整个 stack 链表逐个执行回调任务的tryFire驱动整条回调链a.postComplete();}// result ! null当前下游 this (d) 已经完成// stack ! null当前下游 d 上还有链式回调挂在 stack 链表上if(result!nullstack!null){if(mode0)// 避免同步递归深度爆炸栈溢出所以不调用 postCompletereturnthis;elsepostComplete();}returnnull;}postComplete遍历 stack 中的任务调用 tryFire 执行任务finalvoidpostComplete(){// this - 最开始触发 postComplete 的源 futureCompletableFuture?fthis;Completionh;// h f.stack ! nullf 的回调栈不为空还有回调要执行// f ! this 说明 tryFire 返回了新的下游 future d切换到 d 去处理// 如果 d 栈不为空重新拿 this原始future继续循环while((hf.stack)!null||(f!this(h(fthis).stack)!null)){CompletableFuture?d;Completiont;// CAS 弹出栈顶把 f.stack 从 h 替换成 t h.nextif(f.casStack(h,th.next)){// t ! null代表弹出h之后栈里面还有其他任务if(t!null){if(f!this){// 把刚弹出来的任务 h 重新压回 f 的栈pushStack(h);continue;}// 把h.next置null断开链表h.nextnull;}// 执行回调任务// 这里对 f 重新赋值所以上面的代码需要判断 f ! this// 返回值d不为null代表这个回调产生了新的未完成的 CompletableFuture需要继续处理它的回调栈// 返回null回调执行完毕没有新产生的futuref(dh.tryFire(NESTED))null?this:d;}}}UniApplytryFire判断下游 future 是否已经结束运行或者 uniApply 是否执行失败如果是则返回 null释放属性避免内存泄漏任务运行成功后调用 postFirefinalCompletableFutureVtryFire(intmode){CompletableFutureVd;CompletableFutureTa;// dep 下游 Future就是 thenApply()返回给用户的新 CompletableFuture// src 依赖的那一个上游 CompletableFuture// fn 用户传入的函数Function/Consumer/Runnable// mode 触发模式SYNC-1、ASYNC0、NESTED1// dep null这个回调任务已经执行完毕已经被销毁// uniApply 返回 false 上游 a 还未完成本次不能执行业务回调if((ddep)null||!d.uniApply(asrc,fn,mode0?null:this))returnnull;// 执行到这里说明 dep ! null d.uniApply 返回 true// 切断强引用防止内存泄漏// 对象会被保存在上游 stack 链表如果不清引用会强持有上下游 future 与业务函数即使业务已经结束对象无法回收depnull;srcnull;fnnull;returnd.postFire(a,mode);}总结uniApplyStage 为什么 push 之后立刻调用tryFire(SYNC)push 把任务挂到上游的回调链表但是在 push 完一瞬间上游 Future有可能刚好已经完成了如果不立刻 tryFire要等到别的线程完成上游之后才会触发回调会有延迟tryFire(SYNC)检查上游是否已经完成如果已经完成直接就在当前线程执行这个任务 c如果上游还没完成什么都不做等待上游完成后自动触发 fire。thenApply 调用链路thenApply → uniApplyStage → new UniApply → push(c) → c.tryFire(SYNC) → d.uniApply(src,fn,this) → uniApply执行用户函数complete下游d → dep/src/fn null; → d.postFire(src, SYNC) → 处理上游src的stack → 下游 d 如果有stackmode-1 直接 return d → tryFire 返回 d相关考点CompletableFuture 实现了什么接口CompletableFuture实现FutureCompletionStageFuture传统异步只能阻塞 get 拿结果无法回调CompletionStage定义链式回调接口thenApply /thenAccept/thenCombine…支持非阻塞链式编排supplyAsync / runAsync 区别supplyAsync(Supplier)有返回值runAsync(Runnable)无返回值thenApply / thenAccept / thenRun 区别thenApply有返回值转换结果thenAccept消费结果无返回thenRun不关心上一步结果只执行动作CompletableFuture 内部怎么保存回调采用无锁单向链表栈CAS 操作f.thenApply(fn)生成下游 d把UniApply节点 push 到上游 f 的 stack当 f 执行完成result ! null调用postComplete()遍历 stack 链表逐个执行回调节点tryFire()uniApplyStage 完整流程参数判空新建下游 future d调用d.uniApply(this,f,null)做短路判断如果上游已经完成result!null直接同步执行 fn、完成下游 d返回 true不走压栈如果上游还没完成返回 false → 创建UniApply回调节点调用push(c)压入上游 this 的 stack 链表调用c.tryFire(SYNC)尝试立刻执行返回下游 d。tryFire () 干什么调用上游 future 的uniApply()尝试执行函数执行成功后把引用置空帮助 GC调用下游postFire()modeSYNC(0)同步调用ASYNC(1)异步提交线程池NESTED(-1)postComplete 链表迭代时嵌套调用。postFire () 作用如果上游还有未处理 stack根据 mode 选择cleanStack()清理无效节点或者调用上游postComplete()继续驱动上游链表回调下游自己如果已经完成也触发自己的 postComplete实现链式传播如果回调链中某一步异常后续 thenApply 还会执行吗不会执行。CompletableFuture 回调链中某一步抛出异常后后续的thenApply/thenAccept/thenRun等正常回调会被全部跳过异常状态会沿着链一直向下传播直到遇到能处理异常的方法CompletableFutureStringfutureCompletableFuture.supplyAsync(()-start).thenApply(s-{thrownewRuntimeException(boom);// 第2步异常}).thenApply(s-{System.out.println(不会执行);// 被跳过returns!;}).exceptionally(ex-{System.out.println(捕获: ex.getMessage());// 这里执行returnfallback;});
返回列表