
《Nacos 2.x源码深度解析》专栏目录 一、架构通信篇: 《Nacos 2.x 源码深度解析 (一):架构整体全貌 —— 核心模块划分与版本演进》 《Nacos 2.x 源码深度解析 (二):通信协议迭代 —— HTTP长轮询到gRPC演进》 二、配置中心篇 《Nacos 2.x 源码深度解析 (三):配置中心客户端 —— 启动加载与自动装配》 《Nacos 2.x 源码深度解析 (四):配置中心服务端 —— 事件总线与数据持久化》 《Nacos 2.x 源码深度解析 (五):gRPC 推送链路 —— 配置变更下发与动态刷新》 《Nacos 2.x 源码深度解析 (六):三级缓存体系 —— 降级兜底与故障自愈机制》 三、服务注册发现篇 《Nacos 2.x 源码深度解析 (七):服务注册流程 —— 客户端上报与服务端存储》 《Nacos 2.x 源码深度解析 (八):服务订阅机制 —— 从首次订阅到gRPC双向流变更通知》 四、grpc连接内核篇 《Nacos 2.x 源码深度解析 (九):双向流设计 —— 连接创建复用与销毁》 《Nacos 2.x 源码深度解析 (十):心跳保活策略 —— 断线检测与重连源码》 《Nacos 2.x 源码深度解析 (十一):RPC 请求调度 —— 收发模型与线程池处理》
文章目录
-
- 一、RPC 通信模型总览
-
- 1.1 从连接生命周期到请求调度
- 1.2 三合一通信接口:Requester
- 1.3 请求与响应的序列化协议:protobuf Payload
- 二、客户端请求发送模型
-
- 2.1 三种请求模式的实现与选择
-
- 2.1.1 同步请求:RpcClient.request()
- 2.1.2 异步请求:RpcClient.asyncRequest()
- 2.1.3 Future 模式:RpcClient.requestFuture()
- 2.2 客户端发送模型的设计要点
- 三、服务端请求接收与传输层拦截
-
- 3.1 传输层连接追踪:AddressTransportFilter
- 3.2 连接上下文注入:GrpcConnectionInterceptor
- 3.3 传输层拦截与业务层过滤的分工对比
- 四、服务端请求路由与处理核心:GrpcRequestAcceptor
-
- 4.1 Unary请求的完整处理流程
- 4.2 Handler自动注册机制:RequestHandlerRegistry
- 4.3 Handler过滤链:RequestHandler.handleRequest()
-
- 4.3.1 TpsControlRequestFilter:TPS 限流
- 4.3.2 RemoteRequestAuthFilter:鉴权
- 4.3.3 RemoteParamCheckFilter:参数校验
- 4.4 请求上下文管理
- 五、服务端双线程池架构
-
- 5.1 双线程池的隔离设计
- 5.2 线程池参数配置
- 5.3 线程池的运行时角色
- 六、服务端双向流推送与 ACK 同步
-
- 6.1 服务端推送的两种模式
- 6.2 客户端回执与 ACK 同步器
- 6.3 推送队列反压机制
- 6.4 客户端侧接收服务端推送
- 七、整体链路串联:一次 RPC 请求的完整旅程
- 全文小结
在上一篇文章中,我们深入分析了心跳保活策略,理解了 Nacos 2.x 如何通过 RpcClient.healthCheck() 实现客户端侧的保活探测,通过 reconnect() 实现指数退避重连,以及服务端如何通过 ClientBeatCheckTaskV2 实现两阶段超时降级清理。当连接建立且稳定运行后,所有业务功能——配置查询、服务注册、订阅通知都依赖同一个核心能力:RPC 请求的发送与接收。
但具体而言,Nacos 客户端通过哪三种模式发起 RPC 请求?同步请求的重试与自动切换机制是如何工作的?服务端 gRPC 请求经过哪几个阶段完成从建连到业务处理的完整调度?RequestHandlerRegistry 如何按请求类型自动路由到对应的 Handler?服务端采用怎样的双线程池架构,SDK 与集群请求如何实现物理隔离?业务层的 TPS 限流、参数校验、鉴权三个过滤器的执行顺序和职责分别是什么?服务端通过双向流推送请求时,ACK 回调同步与推送队列反压是如何实现的?本文将沿着一次客户端请求的完整生命周期为主线,从客户端发起、服务端接收、线程池调度、Handler 路由、业务过滤链执行、到响应返回与推送 ACK,串联一条完整的 RPC 请求调度链路。
一、RPC 通信模型总览
1.1 从连接生命周期到请求调度
第 9 篇和第 10 篇完整覆盖了 gRPC 连接的创建、复用、销毁与心跳保活。当连接建立并处于 RUNNING 状态后,请求的发送与接收成为核心关注点。Nacos 2.x 的 RPC 通信基于 gRPC 框架,但并非简单调用 gRPC API。在客户端侧,RpcClient 封装了同步、异步、Future 三种请求模式,内置重试与自动切换服务器;在服务端侧,请求经过传输层拦截器 -> 协议解析 -> Handler 路由 -> 业务过滤链 -> 业务处理五个阶段的完整链路后返回响应;而支撑这一切的底层是精心设计的双线程池架构和双向流推送的 ACK 同步与反压保护机制。
1.2 三合一通信接口:Requester
Nacos 的客户端连接与服务端连接虽然处于不同的模块、承担不同的职责,但它们共同实现了 Requester 接口。这一设计保证了发送模型的统一性——无论是客户端向服务端发送 Unary 请求,还是服务端向客户端推送双向流消息,都遵循同一套契约。
com.alibaba.nacos.api.remote.Requester
public interface Requester {
// 同步请求,阻塞等待直到超时或响应返回
Response request(Request request, long timeoutMills) throws NacosException;
// Future 模式,返回可轮询或阻塞获取的 Future
RequestFuture requestFuture(Request request) throws NacosException;
// 异步回调模式,通过 RequestCallBack 接收响应
void asyncRequest(Request request, RequestCallBack requestCallBack) throws NacosException;
// 关闭连接
void close();
}
该接口定义了三种请求模式:同步(request)、Future(requestFuture)、异步回调(asyncRequest)。客户端模块的 GrpcConnection(common 模块)和服务端模块的 GrpcConnection(core 模块)都实现了此接口,但实现策略不同——客户端的三种模式走 gRPC Unary 调用,服务端的三种模式走双向流推送加 ACK 同步。
1.3 请求与响应的序列化协议:protobuf Payload
所有跨网络的 RPC 请求与响应都通过 protobuf 序列化为统一的 Payload 结构。Payload 是一个二级结构:
- metadata:携带 type(请求类型名,如 "InstanceRequest")、clientIp、headers 等元信息
- body:存放序列化后的具体请求/响应对象字节
message Payload {
Payload.MetaData metadata = 1;
google.protobuf.Any body = 2;
}
message Payload.MetaData {
string type = 1;
string clientIp = 2;
map<string, string> headers = 3;
}
GrpcUtils.convert() 负责将 Request / Response 对象转换为 protobuf Payload,GrpcUtils.parse() 负责反向解析。这一转换发生在每一个 RPC 调用的两端:客户端发送前(GrpcConnection.request() -> GrpcUtils.convert(request))、服务端接收后(GrpcRequestAcceptor.request() -> GrpcUtils.parse(grpcRequest))。
二、客户端请求发送模型
2.1 三种请求模式的实现与选择
RpcClient 作为客户端 SDK 的核心抽象,封装了三种请求发送模式。在分析具体实现之前,先通过下面的时序图建立整体认知:

上面的时序图清晰地展示了三种请求模式的区别:同步模式走完整的 request -> retry -> switchServer 路径;异步模式通过 Guava Futures.addCallback 注册回调,不阻塞调用线程;Future 模式返回 RequestFuture 包装器,调用者可自主选择阻塞获取或轮询。
2.1.1 同步请求:RpcClient.request()
同步请求是使用最频繁的模式。RpcClient.request() 内部实现了完整的重试 + 自动切换机制:
com.alibaba.nacos.common.remote.client.RpcClient#request
public Response request(Request request, long timeoutMills) throws NacosException {
int retryTimes = 0;
Response response;
Throwable exceptionThrow = null;
long start = System.currentTimeMillis();
// 循环重试:默认最多重试 3 次,且在超时时间内持续尝试
while (retryTimes <= rpcClientConfig.retryTimes()
&& (timeoutMills <= 0 || System.currentTimeMillis() < timeoutMills + start)) {
boolean waitReconnect = false;
try {
if (this.currentConnection == null || !isRunning()) {
waitReconnect = true;
throw new NacosException(NacosException.CLIENT_DISCONNECT,
"Client not connected, current status:" + rpcClientStatus.get());
}
// 委托给当前连接的 request 方法发送 gRPC 请求
response = this.currentConnection.request(request, timeoutMills);
if (response == null) {
throw new NacosException(SERVER_ERROR, "Unknown Exception.");
}
if (response instanceof ErrorResponse) {
// 关键:服务端返回 UN_REGISTER 表示连接未注册,触发异步切换服务器
if (response.getErrorCode() == NacosException.UN_REGISTER) {
synchronized (this) {
waitReconnect = true;
if (rpcClientStatus.compareAndSet(RpcClientStatus.RUNNING,
RpcClientStatus.UNHEALTHY)) {
switchServerAsync();
}
}
}
throw new NacosException(response.getErrorCode(), response.getMessage());
}
// 成功则更新上次活跃时间戳并返回
lastActiveTimeStamp = System.currentTimeMillis();
return response;
} catch (Throwable e) {
// 连接断开时等待重连(最多 timeout/3 毫秒)后重试
if (waitReconnect) {
Thread.sleep(Math.min(100, timeoutMills / 3));
}
exceptionThrow = e;
}
retryTimes++;
}
// 所有重试均失败,异步切换服务器
if (rpcClientStatus.compareAndSet(RpcClientStatus.RUNNING, RpcClientStatus.UNHEALTHY)) {
switchServerAsyncOnRequestFail();
}
// 抛出最后一次的异常
if (exceptionThrow != null) {
throw (exceptionThrow instanceof NacosException) ? (NacosException) exceptionThrow
: new NacosException(SERVER_ERROR, exceptionThrow);
} else {
throw new NacosException(SERVER_ERROR, "Request fail, unknown Error");
}
}
该方法的核心设计可以概括为三点:
一是循环重试:内部 while 循环在 retryTimes <= retryTimes()(默认 3 次)且未超时前持续重试,每次失败后短暂休眠(连接断开时休眠 min(100, timeout/3) 毫秒等待重连)。
二是**UN_REGISTER 特殊处理**:当服务端返回 UN_REGISTER 错误码时,说明当前连接已被服务端注销(如服务端重启或心跳超时断开),此时立即将状态置为 UNHEALTHY 并调用 switchServerAsync() 异步切换服务器(不阻塞当前请求的重试循环)。
三是失败兜底:当循环结束后所有重试均失败,调用 switchServerAsyncOnRequestFail() 异步切换,确保连接恢复后后续请求能正常发送。
2.1.2 异步请求:RpcClient.asyncRequest()
异步请求适用于不需要阻塞等待响应结果的场景,如批量请求或非关键路径操作:
com.alibaba.nacos.common.remote.client.RpcClient#asyncRequest
public void asyncRequest(Request request, RequestCallBack callback) throws NacosException {
int retryTimes = 0;
Throwable exceptionToThrow = null;
long start = System.currentTimeMillis();
while (retryTimes <= rpcClientConfig.retryTimes()
&& System.currentTimeMillis() < start + callback.getTimeout()) {
boolean waitReconnect = false;
try {
if (this.currentConnection == null || !isRunning()) {
waitReconnect = true;
throw new NacosException(NacosException.CLIENT_DISCONNECT,
"Client not connected.");
}
// 委托给当前连接,内部通过 Futures.addCallback 注册回调
this.currentConnection.asyncRequest(request, callback);
return; // 异步:发起后立即返回
} catch (Throwable e) {
if (waitReconnect) {
Thread.sleep(Math.min(100, callback.getTimeout() / 3));
}
exceptionToThrow = e;
}
retryTimes++;
}
// 重试耗尽:异步切换服务器
if (rpcClientStatus.compareAndSet(RpcClientStatus.RUNNING, RpcClientStatus.UNHEALTHY)) {
switchServerAsyncOnRequestFail();
}
if (exceptionToThrow != null) {
throw (exceptionToThrow instanceof NacosException) ? (NacosException) exceptionToThrow
: new NacosException(SERVER_ERROR, exceptionToThrow);
}
}
底层 GrpcConnection.asyncRequest() 的实现巧妙利用了 Guava 的 ListenableFuture 和 Futures.addCallback:
com.alibaba.nacos.common.remote.client.grpc.GrpcConnection#asyncRequest
public void asyncRequest(Request request, final RequestCallBack requestCallBack) throws NacosException {
Payload grpcRequest = GrpcUtils.convert(request);
ListenableFuture<Payload> requestFuture = grpcFutureServiceStub.request(grpcRequest);
// 注册成功与失败回调
Futures.addCallback(requestFuture, new FutureCallback<Payload>() {
public void onSuccess(@Nullable Payload grpcResponse) {
Response response = (Response) GrpcUtils.parse(grpcResponse);
if (response != null) {
if (response instanceof ErrorResponse) {
requestCallBack.onException(
new NacosException(response.getErrorCode(), response.getMessage()));
} else {
requestCallBack.onResponse(response);
}
} else {
requestCallBack.onException(
new NacosException(ResponseCode.FAIL.getCode(), "response is null"));
}
}
public void onFailure(Throwable throwable) {
if (throwable instanceof CancellationException) {
requestCallBack.onException(
new TimeoutException("Timeout after "
+ requestCallBack.getTimeout() + " milliseconds."));
} else {
requestCallBack.onException(throwable);
}
}
}, requestCallBack.getExecutor() != null
? requestCallBack.getExecutor() : this.executor);
// 通过 Futures.withTimeout 实现超时控制
Futures.withTimeout(requestFuture, requestCallBack.getTimeout(),
TimeUnit.MILLISECONDS, RpcScheduledExecutor.TIMEOUT_SCHEDULER);
}
这里使用了 Futures.addCallback 注册回调,而回调的执行策略是:优先使用 requestCallBack.getExecutor() 中指定的执行器,若为 null 则使用 GrpcConnection.executor(gRPC 回调线程)。超时控制通过 Futures.withTimeout 配合 TIMEOUT_SCHEDULER 调度器实现。
2.1.3 Future 模式:RpcClient.requestFuture()
Future 模式适用于可并行化的请求场景——调用者发起请求后不立即等待,而是先执行其他任务,之后通过 get() 阻塞获取或 isDone() 轮询:
com.alibaba.nacos.common.remote.client.GrpcConnection#requestFuture
public RequestFuture requestFuture(Request request) throws NacosException {
Payload grpcRequest = GrpcUtils.convert(request);
final ListenableFuture<Payload> requestFuture = grpcFutureServiceStub.request(grpcRequest);
// 返回包装后的 RequestFuture,暴露 get() / isDone() 接口
return new RequestFuture() {
public boolean isDone() {
return requestFuture.isDone();
}
public Response get() throws Exception {
Payload grpcResponse = requestFuture.get();
Response response = (Response) GrpcUtils.parse(grpcResponse);
if (response instanceof ErrorResponse) {
throw new NacosException(response.getErrorCode(), response.getMessage());
}
return response;
}
public Response get(long timeout) throws Exception {
Payload grpcResponse = requestFuture.get(timeout, TimeUnit.MILLISECONDS);
Response response = (Response) GrpcUtils.parse(grpcResponse);
if (response instanceof ErrorResponse) {
throw new NacosException(response.getErrorCode(), response.getMessage());
}
return response;
}
};
}
requestFuture() 同样通过 grpcFutureServiceStub.request(grpcRequest) 发起调用,但区别于 request() 方式——它不阻塞等待结果,而是立即返回一个 RequestFuture 包装器。调用者获得 RequestFuture 后(它在 RpcClient 层也有重试逻辑),再到合适的时机调用 get() 获取结果。
2.2 客户端发送模型的设计要点
综合上述三种模式分析,可以提炼出客户端发送模型的几个关键设计:
一是三种模式的适用场景划分。同步模式(request)用于配置查询、服务注册等需要即时返回结果的场景;异步模式(asyncRequest)适用于批量请求或非关键路径操作,通过回调驱动避免线程阻塞;Future 模式(requestFuture)适用于可并行化的请求场景,调用者可以在发起请求后先执行其他计算,之后再选择阻塞获取或轮询结果。
二是重试与自动切换的策略一致性。三种模式在 RpcClient 层都实现了相同的重试逻辑——循环重试默认 3 次,超时前持续尝试,连接断开时短暂休眠等待重连,重试耗尽后触发 switchServerAsyncOnRequestFail()。UN_REGISTER 错误码被特殊处理为触发 switchServerAsync(),确保连接失效时能快速切换到可用服务器。
三是与 gRPC 底层的关系。GrpcConnection 持有 GrpcFutureServiceStub(用于 Unary 调用)和 StreamObserver(用于双向流),request()、asyncRequest()、requestFuture() 走 Unary 通道,而 sendRequest() / sendResponse() 走双向流通道。两种通道独立存在,各司其职。
三、服务端请求接收与传输层拦截
请求从客户端发出后,最先抵达的是服务端的传输层。在进入业务处理之前,gRPC 请求需要经过两个关键的传输层拦截环节:AddressTransportFilter 负责传输层连接追踪,GrpcConnectionInterceptor 负责将连接上下文注入 gRPC Context。
3.1 传输层连接追踪:AddressTransportFilter
AddressTransportFilter 是 gRPC ServerTransportFilter 的实现,它在 gRPC 传输层连接就绪时触发,提取远程 IP/端口、本地端口等信息,构建 connectionId 并注入到 gRPC Attributes 中:
com.alibaba.nacos.core.remote.grpc.AddressTransportFilter#transportReady
public Attributes transportReady(Attributes transportAttrs) {
InetSocketAddress remoteAddress = (InetSocketAddress) transportAttrs
.get(Grpc.TRANSPORT_ATTR_REMOTE_ADDR);
InetSocketAddress localAddress = (InetSocketAddress) transportAttrs
.get(Grpc.TRANSPORT_ATTR_LOCAL_ADDR);
int remotePort = remoteAddress.getPort();
int localPort = localAddress.getPort();
String remoteIp = remoteAddress.getAddress().getHostAddress();
// 构建 connectionId:timestamp_remoteIp_remotePort
Attributes attrWrapper = transportAttrs.toBuilder()
.set(ATTR_TRANS_KEY_CONN_ID,
System.currentTimeMillis() + "_" + remoteIp + "_" + remotePort)
.set(ATTR_TRANS_KEY_REMOTE_IP, remoteIp)
.set(ATTR_TRANS_KEY_REMOTE_PORT, remotePort)
.set(ATTR_TRANS_KEY_LOCAL_PORT, localPort).build();
Loggers.REMOTE_DIGEST.info("Connection transportReady,connectionId = {} ", connectionId);
return attrWrapper;
}
transportReady() 是所有 gRPC 请求处理之前的第一个钩子。它提取 TCP 连接层的远程地址和本地地址,构建 connectionId(格式为 timestamp_remoteIp_remotePort,形如 1704067200000_192.168.1.10_52800)。这个 connectionId 会通过 gRPC Attributes 链路一直传递到后续的拦截器和业务处理器。
对应的清理方法是 transportTerminated(),当 gRPC 传输层检测到连接断开时自动触发,调用 connectionManager.unregister(connectionId) 清理连接。这层清理作为连接关闭的兜底机制,确保即使业务层未及时释放连接,传输层断开时也能触发清理。
3.2 连接上下文注入:GrpcConnectionInterceptor
GrpcConnectionInterceptor 是 gRPC ServerInterceptor 的实现,在 gRPC Unary 和双向流请求进入业务处理之前,将 AddressTransportFilter 中设置的 Attributes 写入 gRPC Context:
com.alibaba.nacos.core.remote.grpc.GrpcConnectionInterceptor#interceptCall
public <T, S> ServerCall.Listener<T> interceptCall(ServerCall<T, S> call, Metadata headers,
ServerCallHandler<T, S> next) {
// 从 call.getAttributes() 中提取 AddressTransportFilter 写入的连接信息
Context ctx = Context.current()
.withValue(GrpcServerConstants.CONTEXT_KEY_CONN_ID,
call.getAttributes().get(GrpcServerConstants.ATTR_TRANS_KEY_CONN_ID))
.withValue(GrpcServerConstants.CONTEXT_KEY_CONN_REMOTE_IP,
call.getAttributes().get(GrpcServerConstants.ATTR_TRANS_KEY_REMOTE_IP))
.withValue(GrpcServerConstants.CONTEXT_KEY_CONN_REMOTE_PORT,
call.getAttributes().get(GrpcServerConstants.ATTR_TRANS_KEY_REMOTE_PORT))
.withValue(GrpcServerConstants.CONTEXT_KEY_CONN_LOCAL_PORT,
call.getAttributes().get(GrpcServerConstants.ATTR_TRANS_KEY_LOCAL_PORT));
// 双向流请求额外注入 Netty Channel
if (GrpcServerConstants.REQUEST_BI_STREAM_SERVICE_NAME
.equals(call.getMethodDescriptor().getServiceName())) {
Channel internalChannel = getInternalChannel(call);
ctx = ctx.withValue(GrpcServerConstants.CONTEXT_KEY_CHANNEL, internalChannel);
}
return Contexts.interceptCall(ctx, call, headers, next);
}
这里的设计原理是:AddressTransportFilter 在传输层写入的 Attributes 在整个 gRPC Call 期间可用,但 Attributes 是直接与 ServerCall 绑定的,无法直接在方法参数间传递。GrpcConnectionInterceptor 通过 gRPC Context(一个与当前请求线程绑定的作用域变量)将这些信息跨方法传递。后续 GrpcRequestAcceptor.request() 和 GrpcBiStreamRequestAcceptor.requestBiStream() 通过 GrpcServerConstants.CONTEXT_KEY_* 静态常量从 Context 中获取连接上下文。
3.3 传输层拦截与业务层过滤的分工对比
Nacos 服务端的拦截体系分为两个层次:
| 传输层 ServerTransportFilter | AddressTransportFilter | NacosGrpcServerTransportFilter | gRPC 传输层建连/断连 | 追踪连接地址,构建 connectionId |
| 传输层 ServerInterceptor | GrpcConnectionInterceptor | NacosGrpcServerInterceptor | 每个 gRPC 请求入口 | 注入连接上下文到 gRPC Context |
| 业务层 AbstractRequestFilter | TpsControlRequestFilter, RemoteRequestAuthFilter, RemoteParamCheckFilter | SPI(预留) | Handler 内部过滤链 | TPS 限流、鉴权、参数校验 |
传输层拦截发生在 gRPC 框架内部,在所有业务请求处理之前执行,负责连接维度的元信息注入。业务层过滤发生在 Handler 内部,由 RequestHandler.handleRequest() 模板方法统一调度,负责请求维度的业务校验。两者在职责上完全解耦,各司其职。
四、服务端请求路由与处理核心:GrpcRequestAcceptor
请求经过传输层拦截后,到达服务端 Unary 调用的核心枢纽——GrpcRequestAcceptor.request()。这是请求从传输层进入业务层的唯一入口,负责完成请求路由、连接验证、Payload 反序列化、Handler 调度和响应返回的完整流程。
4.1 Unary请求的完整处理流程

上面的时序图清晰地展示了 Unary 请求从接收到响应的完整调度链路。下面深入分析每个步骤的源码实现:
com.alibaba.nacos.core.remote.grpc.GrpcRequestAcceptor#request
public void request(Payload grpcRequest, StreamObserver<Payload> responseObserver) {
String type = grpcRequest.getMetadata().getType();
// 1. 服务端启动检查:未完成启动时拒绝业务请求
if (!ApplicationUtils.isStarted()) {
Payload payloadResponse = GrpcUtils.convert(
ErrorResponse.build(NacosException.INVALID_SERVER_STATUS,
"Server is starting,please try later."));
responseObserver.onNext(payloadResponse);
responseObserver.onCompleted();
return;
}
// 2. ServerCheckRequest 特例处理,不经过 Handler 路由
if (ServerCheckRequest.class.getSimpleName().equals(type)) {
Payload serverCheckResponseP = GrpcUtils.convert(
new ServerCheckResponse(
GrpcServerConstants.CONTEXT_KEY_CONN_ID.get(), true));
responseObserver.onNext(serverCheckResponseP);
responseObserver.onCompleted();
return;
}
// 3. 按请求类型查找对应的 RequestHandler
RequestHandler requestHandler = requestHandlerRegistry.getByRequestType(type);
if (requestHandler == null) {
responseObserver.onNext(GrpcUtils.convert(
ErrorResponse.build(NacosException.NO_HANDLER,
"RequestHandler Not Found")));
responseObserver.onCompleted();
return;
}
// 4. 连接有效性验证
String connectionId = GrpcServerConstants.CONTEXT_KEY_CONN_ID.get();
boolean requestValid = connectionManager.checkValid(connectionId);
if (!requestValid) {
responseObserver.onNext(GrpcUtils.convert(
ErrorResponse.build(NacosException.UN_REGISTER,
"Connection is unregistered.")));
responseObserver.onCompleted();
return;
}
// 5. 反序列化和 Handler 路由
Object parseObj = GrpcUtils.parse(grpcRequest);
if (!(parseObj instanceof Request)) { /* 异常处理 */ }
Request request = (Request) parseObj;
// 6. 构建 RequestMeta 并刷新连接活跃时间
Connection connection = connectionManager.getConnection(connectionId);
RequestMeta requestMeta = new RequestMeta();
requestMeta.setClientIp(connection.getMetaInfo().getClientIp());
requestMeta.setConnectionId(connectionId);
requestMeta.setClientVersion(connection.getMetaInfo().getVersion());
requestMeta.setLabels(connection.getMetaInfo().getLabels());
requestMeta.setAbilityTable(connection.getAbilityTable());
connectionManager.refreshActiveTime(requestMeta.getConnectionId());
// 7. 准备请求上下文
prepareRequestContext(request, requestMeta, connection);
// 8. 执行过滤链 + 业务处理
Response response = requestHandler.handleRequest(request, requestMeta);
Payload payloadResponse = GrpcUtils.convert(response);
// 9. 限流延迟返回:OVER_THRESHOLD 错误延迟 1 秒返回
if (response.getErrorCode() == NacosException.OVER_THRESHOLD) {
RpcScheduledExecutor.CONTROL_SCHEDULER.schedule(() -> {
responseObserver.onNext(payloadResponse);
responseObserver.onCompleted();
}, 1000L, TimeUnit.MILLISECONDS);
} else {
responseObserver.onNext(payloadResponse);
responseObserver.onCompleted();
}
}
Unary 请求的处理流水线可以总结为五个阶段:
阶段一 —— 前置检查:服务端尚未完成启动时拒绝所有业务请求(返回 INVALID_SERVER_STATUS)。ServerCheckRequest 作为特例直接返回 ServerCheckResponse(携带 connectionId),这是客户端 connectToServer() 阶段的第一步握手,不经过 Handler 路由。
阶段二 —— Handler 路由 + 连接验证:通过 requestHandlerRegistry.getByRequestType(type) 按请求类型名查找对应的 RequestHandler。若 Handler 不存在返回 NO_HANDLER。然后检查 connectionManager.checkValid(connectionId),连接未注册则返回 UN_REGISTER(这是客户端触发 switchServerAsync() 的信号)。
阶段三 —— 数据准备:GrpcUtils.parse(grpcRequest) 将 protobuf Payload 反序列化为 Request 对象。然后构建 RequestMeta(携带 connectionId、clientIp、clientVersion、labels、abilityTable),同时调用 refreshActiveTime() 更新连接活跃时间——这是服务端判断客户端是否"存活"的重要依据。
阶段四 —— 业务处理:requestHandler.handleRequest(request, requestMeta) 是过滤链 + 业务逻辑的入口,将在 4.3 节详细分析。
阶段五 —— 后处理与清理:若响应错误码为 OVER_THRESHOLD(TPS 限流),延迟 1 秒返回。最后在 finally 块中调用 RequestContextHolder.removeContext() 清理 ThreadLocal 上下文,防止线程池复用导致上下文污染。
4.2 Handler自动注册机制:RequestHandlerRegistry
RequestHandlerRegistry 是 Handler 路由的核心数据结构。它维护了一个 Map<String, RequestHandler>,key 是请求类型名(如 "InstanceRequest"),value 是对应的 Handler 实例。其自动注册发生在 Spring 容器刷新时:
com.alibaba.nacos.core.remote.RequestHandlerRegistry#onApplicationEvent
public void onApplicationEvent(ContextRefreshedEvent event) {
// 扫描所有 RequestHandler 类型的 Bean
Map<String, RequestHandler> beansOfType =
event.getApplicationContext().getBeansOfType(RequestHandler.class);
Collection<RequestHandler> values = beansOfType.values();
for (RequestHandler requestHandler : values) {
// 遍历继承链找到直接继承 RequestHandler 的子类
Class<?> clazz = requestHandler.getClass();
boolean skip = false;
while (!clazz.getSuperclass().equals(RequestHandler.class)) {
if (clazz.getSuperclass().equals(Object.class)) { skip = true; break; }
clazz = clazz.getSuperclass();
}
if (skip) continue;
// 通过反射提取泛型参数 T 的 SimpleName 作为注册 key
Class tClass = (Class) ((ParameterizedType)
clazz.getGenericSuperclass()).getActualTypeArguments()[0];
registryHandlers.putIfAbsent(tClass.getSimpleName(), requestHandler);
}
}
这里的自动注册机制依赖 Spring 容器:ContextRefreshedEvent 在 Spring 容器刷新完成时触发(通常发生在应用启动或热加载时)。注册流程分为三步:
首先,context.getBeansOfType(RequestHandler.class) 收集容器中所有 RequestHandler 类型的 Bean。
其次,通过反射遍历继承链,找到直接继承 RequestHandler<T, S> 的子类,通过 clazz.getGenericSuperclass() 提取泛型参数 T 的 SimpleName(如 InstanceRequest.class.getSimpleName() 得到 "InstanceRequest")作为 key。
最后,注册到 registryHandlers。同时,若 Handler 标注了 @InvokeSource 注解,还注册到 sourceRegistry,用于 BaseGrpcServer.handleCommonRequest() 中的源合法性校验(判断 SDK 请求是否只路由到 SDK 端口,集群请求是否只路由到集群端口)。若 Handler 的 handle() 方法标注了 @TpsControl,则自动向 TpsControlManager 注册 TPS 点。
4.3 Handler过滤链:RequestHandler.handleRequest()
RequestHandler 是一个抽象模板类,其 handleRequest() 方法定义了标准的业务处理流程——先执行过滤链,再执行业务逻辑:
com.alibaba.nacos.core.remote.RequestHandler#handleRequest
public Response handleRequest(T request, RequestMeta meta) throws NacosException {
// 遍历过滤链,逐个执行 filter()
for (AbstractRequestFilter filter : requestFilters.filters) {
try {
Response filterResult = filter.filter(request, meta, this.getClass());
// 若某个 Filter 返回非成功的 Response,立即短路返回
if (filterResult != null && !filterResult.isSuccess()) {
return filterResult;
}
} catch (Throwable throwable) {
Loggers.REMOTE.error("filter error", throwable);
}
}
// 所有 Filter 通过后,执行业务逻辑
return handle(request, meta);
}
过滤链采用中断模式——TpsControlRequestFilter、RemoteRequestAuthFilter、RemoteParamCheckFilter 三个 Filter 按顺序执行,一旦某个 Filter 返回非 null 且 isSuccess() == false 的 Response,立即中断过滤链并返回该响应。若所有 Filter 返回 null(表示通过),则调用 handle() 执行业务逻辑。
RequestFilters 是一个简单的 ArrayList 集合,Filter 通过 @PostConstruct 注解自动注册:
com.alibaba.nacos.core.remote.AbstractRequestFilter#init
@PostConstruct
public void init() {
requestFilters.registerFilter(this);
}
三个内置 Filter 的执行顺序由其依赖注入的顺序决定(Spring 默认按类名的字母序),在 2.4.3 中的顺序是:
TpsControlRequestFilter —— TPS 限流
RemoteRequestAuthFilter —— 鉴权
RemoteParamCheckFilter —— 参数校验
4.3.1 TpsControlRequestFilter:TPS 限流
com.alibaba.nacos.core.control.remote.TpsControlRequestFilter#filter
protected Response filter(Request request, RequestMeta meta, Class handlerClazz) {
Method method = getHandleMethod(handlerClazz); // 获取 handler 的 handle() 方法
// 检查 handle() 方法是否标注 @TpsControl 且 TPS 控制已启用
if (method.isAnnotationPresent(TpsControl.class)
&& TpsControlConfig.isTpsControlEnabled()) {
TpsControl tpsControl = method.getAnnotation(TpsControl.class);
String pointName = tpsControl.pointName();
// 通过 RemoteTpsCheckRequestParser 解析限流参数
RemoteTpsCheckRequestParser parser =
RemoteTpsCheckRequestParserRegistry.getParser(pointName);
TpsCheckRequest tpsCheckRequest;
if (parser != null) {
tpsCheckRequest = parser.parse(request, meta);
} else {
tpsCheckRequest = new TpsCheckRequest();
}
if (StringUtils.isBlank(tpsCheckRequest.getPointName())) {
tpsCheckRequest.setPointName(pointName);
}
// 执行 TPS 检查
TpsCheckResponse check = tpsControlManager.check(tpsCheckRequest);
if (!check.isSuccess()) {
// 超过限流阈值:通过反射构建默认 Response 实例,设置 OVER_THRESHOLD 错误码
Response response = super.getDefaultResponseInstance(handlerClazz);
response.setErrorInfo(NacosException.OVER_THRESHOLD,
"Tps Flow restricted:" + check.getMessage());
return response;
}
}
return null;
}
该 Filter 首先通过反射获取 handler 的 handle() 方法,检查其是否标注了 @TpsControl 注解。若标注且 TPS 控制已启用,则解析限流参数并调用 tpsControlManager.check() 检查是否超过阈值。超限时通过 getDefaultResponseInstance() 反射创建该 Handler 对应的默认 Response 实例(利用泛型参数 S),设置 OVER_THRESHOLD 错误码返回。
4.3.2 RemoteRequestAuthFilter:鉴权
com.alibaba.nacos.core.auth.RemoteRequestAuthFilter#filter(缩减版)
public Response filter(Request request, RequestMeta meta, Class handlerClazz) {
Method method = getHandleMethod(handlerClazz);
if (method.isAnnotationPresent(Secured.class) && authConfigs.isAuthEnabled()) {
Secured secured = method.getAnnotation(Secured.class);
if (!protocolAuthService.enableAuth(secured)) return null;
// 解析身份和资源
Resource resource = protocolAuthService.parseResource(request, secured);
IdentityContext identityContext = protocolAuthService.parseIdentity(request);
// 验证身份
boolean result = protocolAuthService.validateIdentity(identityContext, resource);
if (!result) throw new AccessException("Validate Identity failed.");
// 验证权限
String action = secured.action().toString();
result = protocolAuthService.validateAuthority(identityContext,
new Permission(resource, action));
if (!result) throw new AccessException("Validate Authority failed.");
}
return null;
}
该 Filter 检查 handle() 方法是否标注了 @Secured 注解且认证已启用。鉴权流程分为两步:先验证身份(Identity),再验证权限(Authority)。验证失败时返回 NO_RIGHT 错误码。鉴权过程中解析的 IdentityContext 和 Resource 被写入 RequestContext,以供后续业务逻辑使用。
4.3.3 RemoteParamCheckFilter:参数校验
该 Filter 检查 Handler 或其 handle() 方法是否标注了 @ExtractorManager.Extractor 注解。若标注,则通过 ExtractorManager 获取对应的参数提取器,提取请求中的参数列表,再由 ParamCheckerManager 调用活动校验器(Active ParamChecker,如客户端 SDK 参数格式校验)执行校验。校验失败时返回 INVALID_PARAM 错误码。
4.4 请求上下文管理
在请求处理的整个过程中,RequestMeta 和 RequestContextHolder 构成了请求的上下文体系。
RequestMeta 携带了请求的元信息:connectionId(连接标识)、clientIp(客户端 IP)、clientVersion(客户端版本号)、labels(标签)、abilityTable(能力表)。这些信息在 GrpcRequestAcceptor.request() 中构建后,贯穿整个处理链路,供鉴权、限流、日志等模块使用。
RequestContextHolder 基于 ThreadLocal 实现,prepareRequestContext() 在请求处理前写入上下文(包括 requestId、userAgent、请求协议、请求目标类名、应用名、远程 IP/端口等),removeContext() 在 finally 块中清理。由于业务 Handler 复用 gRPC 回调线程,ThreadLocal 的清理至关重要,否则线程池复用时将产生上下文污染——上一个请求的上下文信息泄漏到下一个请求中。
五、服务端双线程池架构
5.1 双线程池的隔离设计
Nacos 服务端为 SDK 客户端和集群通信分别配置了独立的 gRPC Server 实例,每个 Server 实例持有自己独立的线程池,实现了请求处理的完全物理隔离:
com.alibaba.nacos.core.remote.grpc.GrpcSdkServer#getRpcExecutor
public ThreadPoolExecutor getRpcExecutor() {
return GlobalExecutor.sdkRpcExecutor;
}
com.alibaba.nacos.core.remote.grpc.GrpcClusterServer#getRpcExecutor
public ThreadPoolExecutor getRpcExecutor() {
if (!GlobalExecutor.clusterRpcExecutor.allowsCoreThreadTimeOut()) {
GlobalExecutor.clusterRpcExecutor.allowCoreThreadTimeOut(true);
}
return GlobalExecutor.clusterRpcExecutor;
}
两个线程池在 GlobalExecutor 中定义:
com.alibaba.nacos.core.utils.GlobalExecutor
// SDK 线程池:处理所有 SDK 客户端的请求
public static final ThreadPoolExecutor sdkRpcExecutor = new ThreadPoolExecutor(
EnvUtil.getAvailableProcessors(RemoteUtils.getRemoteExecutorTimesOfProcessors()),
EnvUtil.getAvailableProcessors(RemoteUtils.getRemoteExecutorTimesOfProcessors()),
60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(RemoteUtils.getRemoteExecutorQueueSize()),
new ThreadFactoryBuilder().daemon(true)
.nameFormat("nacos-grpc-executor-%d").build());
// 集群线程池:处理集群节点间的内部通信
public static final ThreadPoolExecutor clusterRpcExecutor = new ThreadPoolExecutor(
EnvUtil.getAvailableProcessors(RemoteUtils.getRemoteExecutorTimesOfProcessors()),
EnvUtil.getAvailableProcessors(RemoteUtils.getRemoteExecutorTimesOfProcessors()),
60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(RemoteUtils.getRemoteExecutorQueueSize()),
new ThreadFactoryBuilder().daemon(true)
.nameFormat("nacos-cluster-grpc-executor-%d").build());
这两个线程池通过 NettyServerBuilder.executor(getRpcExecutor()) 在 BaseGrpcServer.startServer() 中设置,作为 gRPC Netty Server 的 worker 线程池。这意味着所有通过该 gRPC Server 接入的请求,其 gRPC 框架级别的网络 I/O 和回调执行都由该线程池处理。
5.2 线程池参数配置
线程池的构造参数通过 RemoteUtils 集中管理:
com.alibaba.nacos.core.utils.RemoteUtils
// CPU 倍率默认值:16(1 << 4)
private static final int REMOTE_EXECUTOR_TIMES_OF_PROCESSORS = 1 << 4;
// 队列容量默认值:16384(1 << 14)
private static final int REMOTE_EXECUTOR_QUEUE_SIZE = 1 << 14;
| corePoolSize | availableProcessors × times | 8核 × 16 = 128 | -Dremote.executor.times.of.processors=N |
| maxPoolSize | 同 corePoolSize(固定大小) | 128 | 同上 |
| workQueue | LinkedBlockingQueue | 16384 | -Dremote.executor.queue.size=N |
| keepAliveTime | 固定值 | 60 秒 | 固定 |
| 线程命名(SDK) | nacos-grpc-executor-%d | – | 固定 |
| 线程命名(集群) | nacos-cluster-grpc-executor-%d | – | 固定 |
以 8 核机器为例:corePoolSize = 8 * 16 = 128,队列积压容量 16384,即最多支持 128 个活跃线程同时处理请求,队列中最多积压 16384 个等待任务。
5.3 线程池的运行时角色
需要特别澄清的是,getRpcExecutor() 线程池管理的是 gRPC Netty 的 worker 线程,而非独立的业务线程池。请求从 gRPC 框架到达后,在 Netty EventLoop 线程中依次执行完整的处理链路:

SDK 线程池与集群线程池隔离的核心价值在于:即使 SDK 客户端的请求积压导致 SDK 线程池满载,集群通信线程池不受影响,保证集群节点间的 Distro 同步、心跳复制等核心操作正常执行。这是 Nacos 2.x 服务端高可用的关键设计之一。
六、服务端双向流推送与 ACK 同步
除了处理客户端发起的 Unary 请求外,服务端还需要主动向客户端推送请求——例如配置变更通知、服务端断开连接通知等。这些推送走双向流通道,且需要区分"无等待推送"和"带 ACK 同步推送"两种模式。
6.1 服务端推送的两种模式
com.alibaba.nacos.core.remote.grpc.GrpcConnection
// 模式一:无等待推送
public void sendRequestNoAck(Request request) throws NacosException {
sendQueueBlockCheck(); // 反压检查
// 通过 channel.eventLoop().submit() 确保在 Netty EventLoop 线程中执行
Future<Boolean> executeFuture = this.channel.eventLoop().submit(() -> {
synchronized (streamObserver) {
Payload payload = GrpcUtils.convert(request);
streamObserver.onNext(payload);
return true;
}
});
executeFuture.get(); // 阻塞等待发送完成
}
// 模式二:带 ACK 推送
private DefaultRequestFuture sendRequestInner(Request request, RequestCallBack callBack)
throws NacosException {
// 生成全局唯一 requestId
final String requestId = String.valueOf(PushAckIdGenerator.getNextId());
request.setRequestId(requestId);
// 创建 DefaultRequestFuture,注册回调
DefaultRequestFuture defaultPushFuture = new DefaultRequestFuture(
getMetaInfo().getConnectionId(), requestId, callBack,
() -> RpcAckCallbackSynchronizer.clearFuture(
getMetaInfo().getConnectionId(), requestId));
// 向 RpcAckCallbackSynchronizer 注册待回调的 Future
RpcAckCallbackSynchronizer.syncCallback(
getMetaInfo().getConnectionId(), requestId, defaultPushFuture);
// 调用无等待推送发送
sendRequestNoAck(request);
return defaultPushFuture;
}
无等待推送(sendRequestNoAck):通过 channel.eventLoop().submit() 将发送任务提交到 Netty EventLoop 线程,确保在正确的线程上下文中执行 streamObserver.onNext()。使用 synchronized(streamObserver) 保证并发安全——gRPC 的 StreamObserver#onNext() 不是线程安全的,多线程同时调用可能导致 protobuf 编码时的直接内存泄漏。
带 ACK 推送(sendRequestInner):在无等待推送的基础上增加了 ACK 同步机制。先生成全局唯一的 requestId(通过 PushAckIdGenerator),创建 DefaultRequestFuture 注册到 RpcAckCallbackSynchronizer,然后调用 sendRequestNoAck() 发送。客户端收到请求后通过双向流返回 Response,服务端的 GrpcBiStreamRequestAcceptor.onNext() 识别为 Response 后调用 RpcAckCallbackSynchronizer.ackNotify() 唤醒对应的 Future。
6.2 客户端回执与 ACK 同步器
RpcAckCallbackSynchronizer 是双向流 ACK 同步的核心组件:
com.alibaba.nacos.core.remote.RpcAckCallbackSynchronizer
// 存储结构:ConcurrentLinkedHashMap<connectionId, Map<requestId, DefaultRequestFuture>>
public static final Map<String, Map<String, DefaultRequestFuture>> CALLBACK_CONTEXT =
new ConcurrentLinkedHashMap.Builder<String, Map<String, DefaultRequestFuture>>()
.maximumWeightedCapacity(1000000)
.listener((s, pushCallBack) ->
pushCallBack.entrySet().forEach(entry ->
entry.getValue().setFailResult(new TimeoutException())))
.build();
该组件采用 ConcurrentLinkedHashMap 作为存储结构。外层 key 为 connectionId,内层 key 为 requestId,value 为 DefaultRequestFuture。最大容量 100 万个待回调 Future,当达到容量上限时,最旧的条目被自动逐出并设置 TimeoutException。
关键方法:
-
syncCallback(connectionId, requestId, future):注册待回调的 Future。如果 connectionId 维度的 Map 尚不存在,自动创建(初始容量 128)。
-
ackNotify(connectionId, response):收到客户端回执时唤醒对应的 Future。根据 connectionId 找到对应 Map,按 requestId 移除并设置响应结果。连接过期或请求已超时时,在日志中记录"warn"级别的提示。
-
clearContext(connectionId):连接关闭时清理该连接下所有待回调的 Future。这由 ConnectionManager.unregister() 触发,确保连接断开时所有等待该连接的推送调用不会永久阻塞。
6.3 推送队列反压机制
在 sendRequestNoAck() 的第一步调用了 sendQueueBlockCheck(),这是推送队列的反压保护机制:
com.alibaba.nacos.core.remote.grpc.GrpcConnection#sendQueueBlockCheck
private void sendQueueBlockCheck() {
if (streamObserver instanceof ServerCallStreamObserver) {
// gRPC 内部写缓存阈值固定为 32KB,超过则 isReady() 返回 false
// 见 io.grpc.internal.AbstractStream.TransportState.DEFAULT_ONREADY_THRESHOLD
boolean ready = ((ServerCallStreamObserver<?>) streamObserver).isReady();
if (!ready) {
// 队列不可写:记录 TPS 点 SERVER_PUSH_BLOCK
if (tpsControlManager == null) {
synchronized (GrpcConnection.class.getClass()) {
if (tpsControlManager == null) {
tpsControlManager = ControlManagerCenter.getInstance()
.getTpsControlManager();
tpsControlManager.registerTpsPoint("SERVER_PUSH_BLOCK");
}
}
}
TpsCheckRequest tpsCheckRequest = new TpsCheckRequest("SERVER_PUSH_BLOCK",
this.getMetaInfo().getConnectionId(),
this.getMetaInfo().getClientIp());
tpsControlManager.check(tpsCheckRequest);
getMetaInfo().recordPushQueueBlockTimes();
// 抛异常阻止继续发送
throw new ConnectionBusyException(
"too much bytes on sending queue of this stream.");
} else {
getMetaInfo().clearPushQueueBlockTimes();
}
}
}
反压机制的底层原理是:gRPC 内部有一个写缓存阈值(DEFAULT_ONREADY_THRESHOLD,固定为 32KB),当推送速度超过客户端消费速度时,写缓存积累超过 32KB,ServerCallStreamObserver.isReady() 返回 false。此时服务端不再继续堆积数据,而是记录 TPS 点 SERVER_PUSH_BLOCK 并抛出 ConnectionBusyException,防止服务端推送线程被慢消费者拖垮。当队列恢复可写时,isReady() 重新返回 true,推送恢复正常。
6.4 客户端侧接收服务端推送
在客户端侧,双向流的 StreamObserver.onNext() 回调负责解析服务端推送的消息:
com.alibaba.nacos.common.remote.client.grpc.GrpcClient#bindRequestStream
return streamStub.requestBiStream(new StreamObserver<Payload>() {
public void onNext(Payload payload) {
Object parseBody = GrpcUtils.parse(payload);
final Request request = (Request) parseBody;
if (request != null) {
if (request instanceof SetupAckRequest) {
setupRequestHandler.requestReply(request, null);
return;
}
// 处理服务端推送的请求并返回响应
Response response = handleServerRequest(request);
if (response != null) {
response.setRequestId(request.getRequestId());
sendResponse(response); // 通过双向流返回响应
}
}
}
// onError / onCompleted 处理…
});
客户端 onNext() 的处理逻辑分为两步:首先通过 GrpcUtils.parse(payload) 解析出 Request 或 Response 对象,然后区分两种场景:
-
若为 Request(服务端推送的请求,如配置变更通知):调用 RpcClient.handleServerRequest(),遍历 serverRequestHandlers 列表找到匹配的处理器。该列表通过 registerServerRequestHandler() 由上层业务模块(如 ConfigRpcTransportClient 或 NamingGrpcClientProxy)注册。处理完成后通过 sendResponse(response) 将响应通过双向流返回给服务端,触发服务端的 ackNotify()。
-
若为 Response(请求的 Ack):特殊处理——客户端发起的请求不需要走双向流 Ack(走 Unary 响应),这里收到 Response 是服务端主动推送请求时客户端返回的 Ack,但这是服务端层面的处理(服务端 GrpcBiStreamRequestAcceptor.onNext() 中处理 Response 类型)。
实际上,客户端侧的 onNext() 主要在处理服务端推送的 Request,而服务端侧的 GrpcBiStreamRequestAcceptor.onNext() 则在处理客户端返回的 Response(Ack)。双方通过双向流交换数据,但视角不同。
七、整体链路串联:一次 RPC 请求的完整旅程
将前文各节的内容串联起来,一条完整的 RPC 请求调度链路如下:

关键设计思想:Nacos 2.x 的 RPC 请求调度采用了分层解耦与职责分离的设计哲学。传输层(AddressTransportFilter / GrpcConnectionInterceptor)与业务层(AbstractRequestFilter 过滤链)各自专注本层职责;双线程池在 GrpcSdkServer 与 GrpcClusterServer 层面实现物理隔离;RequestHandlerRegistry 的自动注册机制将 Handler 的发现与路由完全交给 Spring 容器管理;RpcAckCallbackSynchronizer + 反压机制保证了双向流推送的可靠性与稳定性。这一整套调度体系在保障高性能的同时,兼顾了可扩展性与运维友好性。
全文小结
本文聚焦 RPC 请求调度,从客户端发送模型和服务端调度链路两个维度,深入分析了 Nacos 2.x 基于 gRPC 框架的完整 RPC 通信架构。
在客户端三模式请求模型方面,RpcClient 封装了同步、异步、Future 三种请求发送模式,同步请求内置重试与 UN_REGISTER 自动切换机制,三种模式底层均通过 GrpcConnection 的 gRPC FutureStub 完成 Unary 调用。在服务端五阶段调度链路方面,请求经过传输层拦截(AddressTransportFilter + GrpcConnectionInterceptor)→ 源合法性检查(handleCommonRequest)→ 核心 Acceptor 调度(GrpcRequestAcceptor.request)→ 业务过滤链(TPS 限流 → 鉴权 → 参数校验)→ 业务 Handler 执行,五个阶段职责清晰。双线程池隔离架构通过 GrpcSdkServer 与 GrpcClusterServer 各自持有独立线程池(sdkRpcExecutor / clusterRpcExecutor),线程池大小公式为 CPU 核数 × 16,队列容量 16384,SDK 与集群请求物理隔离互不影响。Handler 自动注册机制通过 RequestHandlerRegistry 在 Spring 容器刷新时自动扫描所有 RequestHandler Bean,通过泛型参数 T 的 SimpleName 建立请求类型到 Handler 的映射。双向流推送与 ACK 同步方面,服务端通过 sendRequestNoAck(无等待)与 sendRequestInner(带 ACK)两种模式推送请求,RpcAckCallbackSynchronizer 基于 ConcurrentLinkedHashMap 管理待回调 Future(最大 100 万),sendQueueBlockCheck 通过 gRPC isReady() 实现推送队列反压保护。
原创不易,如果本文对您有帮助,带来了些许灵感或启发,烦请动动小手点赞、关注、转发、收藏。这是作者持续更新的动力源泉,衷心感谢您的支持。我会尽量在工作之余,为大家带来更高品质的内容,努力保持周更。




