Jaws 系列第 6 篇。前情提要:《删掉 gRPC 依赖后,我用 2400 行打通了 gRPC 生态》、《HTTP/2 传输进化史》。
本文代码全部出自 javahongxi/jaws,commit 可查。
一、一个让人不舒服的 join()
故事要从 wire 模块的一次重构说起。
8 月底我给 jaws-wire 加非阻塞分发时,WireCallDispatcher里有一段代码让我越看越难受:
CompletableFuture<Object>future=messageHandler.handleAsync(jawsRequest);// 旧版:Responseresult=future.join();// 业务线程就这么被占着等handleDispatchResult(result,ctx,serverHandler);这是 gRPC 线格式的 Provider 管线分发路径。handleAsync明明返回的是CompletableFuture——异步语义都到位了——下游却一个join()把结果等回来。业务方法快还好,一旦下游是慢调用,每个在途请求都占着一个业务线程干等。线程数就是并发的天花板,这等于把异步框架用出了同步的损耗。
改成什么样?先把答案放这里:
CompletableFuture<Object>future=messageHandler.handleAsync(jawsRequest);if(future.isDone()){// 快路径:业务方法同步返回时直接内联处理,不切线程handleDispatchResult(future.join(),ctx,serverHandler);}else{// 慢路径:挂回调,业务线程立刻释放future.whenComplete((result,throwable)->{handleDispatchResult(result,ctx,serverHandler);});}isDone()快路径是个值得停下来品一下的细节:大多数业务方法其实跑得飞快,future 返回时已经完成了,这时候再挂回调纯属浪费一次线程切换。所以先探一下,已完成就内联join()(此时 join 不阻塞),没完成才走whenComplete。一行判断,换掉一次无谓的调度开销。
改完 wire 这段,我以为这事就结束了。直到我顺着调用链往回看了一眼客户端——才发现真正的题目有多大。
二、顺着调用链往上爬
RPC 框架的调用链,从消费端到服务端大致是这样一条河:
业务代码 → 代理层 → Filter 链 → Cluster/LB → Reference → 序列化 → 传输层 ↓ (网络) 传输层 ← 序列化 ← Provider 管线 ← Filter 链 ← 业务实现wire 的join()只是河下游的一个洞。我从洞往上游走,一路上看到的是这样的景象:
Filter 契约是同步的。jaws 的Filter接口长这样:
publicinterfaceFilter{Responsefilter(Caller<?>caller,Requestrequest);}返回Response意味着什么?意味着 Filter 必须拿到最终结果才能返回。而 Filter 的下一跳caller.call(request)走到传输层,传输层得等网络回包——于是每个 Filter 都被迫阻塞等一次完整的 RPC。更别提链式组装了:三个 Filter 串起来,就是三次「等到天荒地老」的串行阻塞。
TracingFilter 已经在偷偷变形。最典型的是链路追踪 Filter,改造前的消费端逻辑:
// 改造前(同步契约下的无奈写法)privateResponsehandleConsumer(...){Spanspan=t.nextSpan().name(spanName).start();try(Tracer.SpanInScopescope=t.withSpan(span)){p.inject(span.context(),request,Request::setAttachment);Responseresponse=caller.call(request);// 阻塞等整个 RPCif(response.getException()!=null){span.error(response.getException());}returnresponse;}catch(Exceptione){span.error(e);throwe;}finally{span.end();}}这段代码有个隐藏的问题:span.end()在 finally 里,时机是「RPC 结束」——没错,但前提是线程一直停在这里等。span 的生命周期和线程的阻塞生命周期被迫绑死了。想做到「请求发出后线程就走,span 等回包时再关」?在同步契约下根本写不出来。
DefaultResponseFuture 是一只自研怪兽。客户端等回包的核心类,改造前 252 行,自带状态机:
publicclassDefaultResponseFutureimplementsResponseFuture{protectedvolatileFutureStatestate=FutureState.DOING;// 自研状态枚举protectedvolatileObjectresult;protectedvolatileExceptionexception;protectedvolatileList<FutureListener>listeners;// 自研监听器publicObjectgetValue(){synchronized(this){if(!isDoing()){returngetValueOrThrow();}longwaitTime=timeout-(System.currentTimeMillis()-createTime);if(waitTime>0){for(;;){try{wait(waitTime);// 手写 wait/notify}catch(InterruptedExceptionignore){Thread.currentThread().interrupt();}if(!isDoing())break;// 重新计算剩余超时,继续 wait …}}if(isDoing())cancelOnTimeout();returngetValueOrThrow();}}}手写 wait/notify 循环、手写超时补偿、手写 listener 通知、FutureState三态枚举、FutureListener回调接口……这些 JDK 在CompletableFuture里全部都有,而且经过十几年生产环境锤炼。我维护这只怪兽的每一行,都在重新发明 JDK 已经解决的问题。
服务端管线也在同步裸奔。顺着河再看服务端一侧:MessageHandler是传输层和业务之间的桥,它在旧版里的签名是同步的Response handle(Object message)。Netty 的NettyChannelHandler收到请求后,得等handle返回才能把响应写回 channel——事件循环线程(event loop)被业务调用整段占住。而事件循环是 Netty 的命根子:一个 event loop 管着几十条连接,它被卡 100ms,这几十条连接上的所有请求全部顺延。这不是「慢一点」的问题,是所有连接互相拖累的放大器。
到这里,问题已经从「wire 有个 join 不顺眼」升级成了:框架名字叫 JAWS(Java Async Wire Service),异步却只做到了传输层往上一点点——Filter 拦在中间,Future 是同步内核,整条链是异步的躯干拖着同步的四肢。
三、决策:契约先行,从 Filter 开刀
改造顺序是这次最有意思的决策点。可选路径有三条:
- 先改传输层:把 Netty/http2 内部全异步化,Filter 契约不动——治标,join 只是藏得更深
- 先改 Filter 契约:
filter()返回CompletableFuture<Response>,逼着整条链跟着改 - 先改 Future:DefaultResponseFuture 换成 CompletableFuture,上层不动
我选了 2,契约先行。理由:契约是骨架,骨架定了肉才知道往哪长。Filter.filter()的返回类型一变,编译器会把所有「不改就会坏」的地方精确地列出来——Filter 实现、包装器、调用方,一个都逃不掉。这比人肉排查安全得多。
新的Filter接口:
@SpipublicinterfaceFilter{CompletableFuture<Response>filter(Caller<?>caller,Requestrequest);}配套在Caller( jaws 里消费端 Reference 和服务端 Provider 的共同抽象,对标 Dubbo 的 Invoker)上加默认方法:
publicinterfaceCaller<T>extendsEndpoint{Responsecall(Requestrequest);defaultCompletableFuture<Response>callAsync(Requestrequest){returnCompletableFuture.completedFuture(call(request));}}callAsync()做成 default 方法是个关键的兼容设计:新契约是「邀请」而不是「强拆」。老实现不改照样编译通过(默认桥接到同步 call),新的实现可以按自己的节奏迁到异步。整个迁移过程中测试始终是绿的。
TracingFilter 的蜕变
契约一换,之前写不出来的代码自然就长出来了:
// 改造后:span 生命周期不再绑死线程privateCompletableFuture<Response>handleConsumer(...){Spanspan=t.nextSpan().name(spanName).start();try(Tracer.SpanInScopescope=t.withSpan(span)){p.inject(span.context(),request,Request::setAttachment);}returncaller.callAsync(request).whenComplete((response,throwable)->{if(throwable!=null){span.error(throwable);}elseif(response!=null&&response.getThrowable()!=null){span.error(response.getThrowable());}span.end();// 回包时关 span,线程早就走了});}对比一下两版的本质差异:注入 trace 上下文这种「出发前」的动作在调用线程完成;span.end()挪进whenComplete,变成「到达后」的动作,由完成 future 的那个线程执行。span 的生命周期终于和 RPC 的生命周期对齐,而不是和某个线程的阻塞周期对齐。这就是异步契约的价值——让代码的形状贴住事物的真实形状。
AccessLog、TokenAuth、Metrics 四个内置 Filter 全部照此迁移,全部改成callAsync().whenComplete()的非阻塞后处理模式。
FilterProviderWrapper:一个 join 也不留
Filter 链的组装节点是FilterProviderWrapper,它的双轨实现最能说明这次升级的彻底性:
@OverridepublicResponsecall(Requestrequest){if(isFilterDisabled(request.getInterfaceName())){returnoriginal.call(request);}returnfilter.filter(original,request).join();// 同步门面,内部一次 join}@OverridepublicCompletableFuture<Response>callAsync(Requestrequest){if(isFilterDisabled(request.getInterfaceName())){returnoriginal.callAsync(request);}returnfilter.filter(original,request);// 原生异步,零 join}同步入口call()保留给确实需要阻塞语义的调用方(内部一次 join 封装),异步入口callAsync()则全程无 join 直通。框架内部全部走callAsync,同步门面只留给用户边界。
四、把 252 行怪兽换成一行 extends
Filter 契约改完,轮到那只怪兽了。
DefaultResponseFuture的改造方向其实早在选型时就定了:JDK 的CompletableFuture已经把状态机、监听器、超时、组合算子全部做好,自研的意义只剩「历史包袱」。于是改造后的全部声明是——
publicclassDefaultResponseFutureextendsCompletableFuture<Response>implementsResponseFuture{privatefinalRequestrequest;privatefinalinttimeout;@OverridepublicvoidonSuccess(Responseresponse){complete(response);}@OverridepublicvoidonFailure(Responseresponse){completeExceptionally(response.getThrowable());}}252 行瘦身到 128 行——剩余部分基本是 Javadoc 和给老 API 用的异常转换。删掉的东西列一下,全是 JDK 替我保管的:
FutureState三态枚举 →CompletableFuture内部状态机FutureListener+ 手写通知循环 →whenComplete/thenApply- 手写 wait/notify 超时循环 → 传输层每请求定时器 +
completeExceptionally - 手写的
getRawValue()/getThrowable()状态判读 →getNow(null)/isCompletedExceptionally()组合
保留ResponseFuture门面接口是为了 API 稳定(getValue()/getTimeout()/getRequestId()这些老签名还在),但内核已经完全是 CompletableFuture 了。阻塞式用户代码getValue()委托给get(),异步用户代码直接whenComplete,同一只 future 两种活法。
有个容易被忽略的收益藏在 commit message 里:「消除双 Future 分配开销和 monitor lock 成本」。旧实现里 Reference 层为了拿到 CompletableFuture 语义,要再 new 一个 CompletableFuture 把 ResponseFuture 桥接过去——每个请求两只 future、一次锁。现在DefaultResponseFuture本身就是 CompletableFuture,AbstractReference 里直接whenComplete链上去,一只 future 贯穿到底:
// AbstractReference.callAsync:网络层完成 future 时顺路做统计if(responseinstanceofCompletableFuture<?>cf){returncf.whenComplete((r,t)->{decrActiveCount(response);if(t==null){longelapsed=System.nanoTime()-startTime;succeededElapsed.addAndGet(elapsed);succeededCount.incrementAndGet();}}).thenApply(r->(Response)r);}调用统计(activeCount、成功耗时累计)本来要靠额外回调挂钩,现在就是 future 链上的一环。基础设施统一之后,横切逻辑的挂载方式自然就变优雅了——这是我认为比「少 300 行」更值钱的部分。
五、收尾:语义对齐与三个细节
主战役之外,这次还顺手做了三件值得单独一说的事。
异常语义对齐。getException全链路改名getThrowable,类型从Exception放宽到Throwable。别小看这个放宽:序列化栈溢出抛的是StackOverflowError,同步业务代码抛Error子类的场景也不罕见,老契约里这些直接被类型系统拒收,只能丢信息。改名加放宽,是一起还的旧债。
服务端管线 handleAsync。前面埋的伏笔在这里兑现:MessageHandler的签名从同步handle()换成handleAsync()返回CompletableFuture<Object>,由AbstractRequestHandler统一实现——Provider 查找、方法解析这些前置逻辑照旧在事件循环完成(纯内存操作,快且安全),真正的业务调用交给doHandleAsync。NettyChannelHandler收到请求后改成挂whenComplete回调,回包到达时再把响应写回 channel。事件循环线程从此只做「接请求、写响应」两件事,中间漫长的业务执行与它无关。
熔断计数修正。WireClient 的错误熔断要对齐 Http2Client 的whenComplete语义——区分「业务异常」和「框架错误」:业务抛JawsBizException说明服务本身活着,只是这次调用失败,不计入熔断;网络错误、超时这类框架级错误才incrErrorCount(),成功后resetErrorCount()。区分的意义在于避免业务侧正常的异常流量把熔断器误打开——框架挂了才该熔断,业务出错不该。
用户侧异步 API 落地。契约改完了,得让用户用得上。DemoService接口补了helloAsync(String 返回)和getUserAsync(POJO 返回)两个异步方法,NettyConsumer 里演示了从基础回调到并发组合的完整用法:
// 基础异步:提交后线程立刻返回,回调里拿结果CompletableFuture<String>asyncHello=demoService.helloAsync("async-lily");asyncHello.thenAccept(result->System.out.println("callback => "+result));// thenApply 链式变换demoService.helloAsync("chain-demo").thenApply(String::toUpperCase).thenAccept(s->System.out.println("chained => "+s));// allOf 并发组合:两个调用齐了再汇总CompletableFuture<String>f1=demoService.helloAsync("user-A");CompletableFuture<String>f2=demoService.helloAsync("user-B");CompletableFuture.allOf(f1,f2).thenRun(()->System.out.println("combined => ["+f1.join()+", "+f2.join()+"]"));六、W 落定:Async Wire Service 名实对齐
回头看这次改造的完整地图:
| 层 | 改造前 | 改造后 |
|---|---|---|
| Filter 契约 | Response filter(...)同步拦截 | CompletableFuture<Response> filter(...)异步管道 |
| Caller 抽象 | 只有call() | 新增callAsync()default 桥接 |
| 客户端 Future | 252 行自研 wait/notify 状态机 | extends CompletableFuture,128 行 |
| 传输层分发 | future.join()占线程 | isDone()快路径 +whenComplete回调 |
| 服务端管线 | 同步handle() | handleAsync()返回 CompletableFuture |
| 流式 | wire/http2 各自实现 | 共享StreamPublisher(Flow.Publisher 基座) |
JAWS 这个名字立项目时就起好了——JavaAsyncWireService。但说实话,改造之前那个 A 是有水分的:异步只在传输层往上一点点,Filter 拦腰一截同步,客户端内核是 wait/notify。这次全链贯通之后,从业务代理、Filter、Cluster、Reference 到传输层再到 wire 分发,CompletableFuture 一竿子插到底,A 字才算名副其实。
复盘这次改造,有三条经验值得留给读者:
- 契约先于实现。改返回类型这种「编译器帮你找全改动点」的路径,永远比人肉排查安全。同步入口保留为门面、异步作为原生路径的双轨设计,让迁移过程没有一天是红的。
- 删自研代码之前,先确认标准库真的覆盖了你的场景。CompletableFuture 覆盖了状态机、监听、超时、组合;而 jaws 特有的每请求超时归传输层定时器管理——边界划清了,删起来才不心虚。
- 顺着调用链找下一个洞。wire 的一个 join() 只是症状,从症状出发往上游走完整条链,才能看到问题的全貌。单点优化容易,全链对齐难,难就难在要舍得把「已经能跑」的代码推翻。
jaws 现在核心四模块约 2.3 万行,README 里那句「可以从头读到尾」依然成立——而且这次重构之后,「读」的体验更好了:你顺着一条callAsync从 Filter 走到传输层,看到的将是同一种异步语言,而不是三种方言的翻译现场。
下一篇候选是自适应负载均衡的实现剖析(power of two choices 在 RPC 里的落地),感兴趣的可以先去仓库读AdaptiveLoadBalance,我们下篇见。
项目地址:github.com/javahongxi/jaws — 核心约 2.3 万行、可从头读到尾的轻量级 RPC 框架。三传输(Netty 二进制 / HTTP/2 / gRPC 线格式),gRPC 互通,Server Streaming,自适应负载均衡,实测 10 万 QPS。