Nacos 远程连接生命周期规范深度解析:gRPC 连接管理、Push/Ack 与运行时踢除机制
【免费下载链接】nacosan easy-to-use dynamic service discovery, configuration and service management platform for building AI cloud native applications.项目地址: https://gitcode.com/GitHub_Trending/na/nacos
本文基于 Nacos 仓库 foundation-remote-connection-spec.md 规范文档,结合
core模块远程层源码,系统讲解 Nacos 服务端远程连接的基础生命周期。你将掌握 Nacos 如何抽象Connection模型、通过双向流完成连接 Setup 与注册、处理 Unary 请求、维护服务端到客户端的 Push 与 Ack 管道,以及如何通过 Active Detection 与运行时踢除保护连接资源——这套机制是 Config、Naming、AI 等所有领域能力在 gRPC 传输层之上的共同底座。
1. 范围:远程连接层负责什么
远程连接层位于领域逻辑与 gRPC 传输之间,负责传输连接状态、请求上下文、推送与 Ack 管道,以及连接保护。它是 基础能力规范 中远程连接部分的展开,但不负责领域请求语义——详细请求上下文与 filter 规则由 请求过滤与运行时上下文规范 定义。
1.1 负责的能力
- SDK 与 cluster 两类来源的 gRPC transport 启动;
- connection id 创建与 transport 地址捕获;
- 连接 setup、注册、拒绝与注销;
- 连接元数据与能力表传递;
- 请求元数据与请求上下文创建;
- 服务端到客户端 push 以及 ack 匹配;
- active time 刷新、过期连接检测与运行时踢除;
- 面向领域运行时状态的连接 listener 回调。
1.2 明确不负责的部分
- Config、Naming、AI 或 auth payload 的语义;
- 公开 HTTP API 行为;
- 持久领域状态;
- SDK 重试、failover 或 server-list 选择行为(见 客户端连接与故障切换规范);
- AP/CP 一致性 membership。
公开 gRPC 包裹与 JSON payload 策略由 gRPC API 规范 定义,本文只讨论承载这些请求的服务端连接生命周期。
2. Connection 模型:四个核心抽象
Connection是服务端对 live remote connection 的抽象,GrpcConnection是默认 gRPC 实现。规范文档定义了四个协作模型:
| 模型 | 职责 |
|---|---|
Connection | 暴露连接元数据、labels、能力表、trace 标记、同步请求、异步请求、no-ack push、close 和连通性检查。 |
ConnectionMeta | 存储 connection id、source、client ip、remote endpoint、local port、version、app、namespace、labels、TLS 标记、创建时间、最后活跃时间和 push queue block 时间戳。 |
ConnectionManager | 拥有内存连接注册表、按 client ip 统计、连接 control 检查、listener 回调、active time 刷新和运行时踢除。 |
ClientConnectionEventListener | 允许领域在连接注册或注销时挂载或释放运行时状态。 |
在源码中,这四个模型分别对应 Connection.java、ConnectionMeta.java、ConnectionManager.java 与 ClientConnectionEventListener.java。其中ConnectionManager内部维护两张核心数据结构(见ConnectionManager.java第 61-63 行):
private Map<String, AtomicInteger> connectionForClientIp = new ConcurrentHashMap<>(16); Map<String, Connection> connections = new ConcurrentHashMap<>();connections:以 connection id 为键的内存连接注册表,提供register/unregister/getConnection/checkValid等操作;connectionForClientIp:按 client ip 统计的连接计数,供连接 control 与监控使用。
2.1 连接身份规则
connectionId在 gRPC transport ready 时由时间戳、remote ip 和 remote port生成;- 普通 unary 请求被接受前,连接必须完成
ConnectionSetupRequest处理; - labels 必须包含 source 信息,尤其是 SDK 或 cluster source;
- cluster-source 连接是内部连接,跳过 SDK 连接数限制检查;
- 能力表在 setup 时协商,并通过
RequestMeta传递给 request handler; - 已接受的 unary 请求和 push ack 必须刷新
lastActiveTime。
connection id 的实际生成逻辑位于 AddressTransportFilter.java 的transportReady回调中,格式为时间戳_远端IP_远端端口:
.set(ATTR_TRANS_KEY_CONN_ID, System.currentTimeMillis() + "_" + remoteIp + "_" + remotePort)同时该 filter 还会把 remote ip、remote port、local port 一并写入 transport attributes,随后由GrpcConnectionInterceptor注入 gRPC context(见下文第 3 节)。cluster-source 跳过连接数限制的实现在ConnectionManager.checkLimit()(第 133-146 行):当connection.getMetaInfo().isClusterSource()为 true 时直接返回 false(不限制),否则通过ControlManagerCenter的连接 control manager 执行ConnectionCheckRequest检查。
3. gRPC Server 表面:两类来源,一个基类
Nacos 暴露两类 gRPC server source:
| 来源 | Server | 端口偏移 | 目的 |
|---|---|---|---|
| SDK | GrpcSdkServer | SDK gRPC port offset | 客户端到服务端请求和服务端 push。 |
| Cluster | GrpcClusterServer | Cluster gRPC port offset | 服务端间请求。 |
两个实现都继承BaseGrpcServer,对应源码为 GrpcSdkServer.java、GrpcClusterServer.java 与 BaseGrpcServer.java。GrpcSdkServer通过rpcPortOffset()返回Constants.SDK_GRPC_PORT_DEFAULT_OFFSET,集群服务端使用独立的 cluster port offset,两者可并行监听。
3.1 两类 server 的共同表面
BaseGrpcServer.addServices()(第 215-259 行)注册了三个关键构件:
- unary
Request/request:普通 request/response 调用,最终路由到GrpcRequestAcceptor.request(); - bidirectional
BiRequestStream/requestBiStream:用于 setup、server push 和 push ack,路由到GrpcBiStreamRequestAcceptor; AddressTransportFilter:捕获 remote/local 地址并生成 connection id,在 transport terminated 时注销连接;GrpcConnectionInterceptor:把 connection id 与地址属性放入 gRPC context;- 每个 source 专属的 interceptor、transport filter、keepalive 配置与 inbound-message-size 配置。
GrpcConnectionInterceptor(GrpcConnectionInterceptor.java)从ServerCall的 attributes 中读取 connection id、remote ip、remote port、local port,写入Context;对于双向流调用还会把底层 NettyChannel一并放入 context,供后续直接读写 stream。
AddressTransportFilter(AddressTransportFilter.java)实现了两个生命周期回调:
transportReady:捕获地址、生成 connection id 并写回 transport attributes;transportTerminated:取出 connection id 并调用connectionManager.unregister(connectionId)完成自动注销。
3.2 服务端 keepalive:传输层保护
服务端 keepalive 配置属于传输层保护。BaseGrpcServer.startServer()通过NettyServerBuilder设置keepAliveTime、keepAliveTimeout、permitKeepAliveTime与maxInboundMessageSize,默认值集中在GrpcServerConstants.GrpcConfig中;GrpcSdkServer允许通过 SDK 专属配置项覆盖 keepalive 时间、超时、permit 时间与最大入站消息大小(见GrpcSdkServer.java第 59-108 行)。
规范文档特别强调:领域模块必须消费连接生命周期事件,而不是自行实现 gRPC 传输层心跳。换言之,健康检测与心跳的统一入口在远程连接层,领域只关心业务事件。
4. Setup 与注册:双向流上的握手
连接 setup 通过 bidirectional stream 完成,规范定义 8 步流程:
- transport 层在 gRPC transport ready 时创建连接属性;
- 客户端在 bidirectional stream 上发送
ConnectionSetupRequest; - 服务端根据 payload metadata、setup labels、client version、namespace、source、app name、local port 和 TLS 状态构建
ConnectionMeta; - 服务端通过
ConnectionGeneratorService创建Connection(源码:ConnectionGeneratorService.java); - setup 中存在能力表时,将能力表存入 connection;
- 服务端仍处于 starting 状态时拒绝 SDK 连接;
ConnectionManager.register校验连接状态、应用连接 control 规则、记录 trace、保存连接、递增 per-ip 计数,并通知 listener;- 如果使用了能力协商,服务端发送带当前服务端能力表的
SetupAckRequest。
注册失败必须关闭连接,且不得留下绑定到被拒绝 connection id 的领域运行时状态。
4.1 register 的源码实现
ConnectionManager.register()(ConnectionManager.java)是第 7 步的具体实现:
public synchronized boolean register(String connectionId, Connection connection) { if (connection.isConnected()) { String clientIp = connection.getMetaInfo().clientIp; if (connections.containsKey(connectionId)) { return true; } if (checkLimit(connection)) { return false; } if (traced(clientIp)) { connection.setTraced(true); } connections.put(connectionId, connection); connectionForClientIp.computeIfAbsent(clientIp, k -> new AtomicInteger(0)).getAndIncrement(); clientConnectionEventListenerRegistry.notifyClientConnected(connection); ... } return false; }要点:
- 只有在
connection.isConnected()为 true 时才允许注册; - 通过
checkLimit应用连接 control 规则(含 cluster source 豁免); - 对被 trace 的 client ip 打开 trace 标记,便于在 GrpcRequestAcceptor.java 中打印 payload 详情;
- 注册成功后立刻通过
ClientConnectionEventListenerRegistry.notifyClientConnected通知所有 listener,让领域挂载运行时状态。
5. Unary 请求处理:GrpcRequestAcceptor 的完整链路
Unary 请求由GrpcRequestAcceptor处理(GrpcRequestAcceptor.java)。规范定义的请求处理规则,在源码中逐条可验证:
- 服务端 starting 时拒绝请求,server check 除外:第 92-103 行,当
!ApplicationUtils.isStarted()时返回INVALID_SERVER_STATUS错误; ServerCheckRequest直接处理并返回 connection id:第 106-115 行,直接构造ServerCheckResponse(CONTEXT_KEY_CONN_ID.get(), true)返回,不经过 handler 与注册表;- request handler 按请求类 simple name 解析:第 117 行
requestHandlerRegistry.getByRequestType(type),未找到 handler 返回NO_HANDLER错误; - gRPC context 中的 connection id 必须已经注册:第 134-149 行
connectionManager.checkValid(connectionId),未注册返回UN_REGISTER错误; - malformed、unknown 或 non-request payload 返回错误响应:第 151-197 行依次处理解析异常、解析为空、非
Request类型三种情况; RequestMeta必须包含来自注册连接的 client ip、connection id、client version、labels 和能力表:第 203-208 行从connection.getMetaInfo()与connection.getAbilityTable()构建;RequestContext必须填充 protocol、request target、app、remote endpoint 和 source ip:第 245-263 行prepareRequestContext完成填充(GRPC_PROTOCOL、请求类 simple name、app、remoteIp/remotePort、sourceIp);- request filter 在领域 handler 之前执行:filter 链通过 RequestFilters.java 与 AbstractRequestFilter.java 承载,规则见 请求过滤与运行时上下文规范;
- 已接受的 unary 请求会在进入领域处理前刷新连接 active time:第 209 行
connectionManager.refreshActiveTime(requestMeta.getConnectionId())。
Request handler 拥有 payload 语义。远程层只负责路由解析、元数据、上下文、filter 与响应转换。对
OVER_THRESHOLD响应,源码会延迟 1 秒返回(第 214-219 行),用于控制面限流场景。
6. Push 与 Ack:服务端主动推送的可靠性管道
服务端到客户端 push 使用已注册的Connection。核心类为 RpcPushService.java(发起 push)与 RpcAckCallbackSynchronizer.java(ack future 匹配)。Push 规则:
sendRequestNoAck将 NacosRequest转换为 gRPC payload,并写入 bidirectional stream;- 由于
StreamObserver.onNext不是线程安全的,stream 写入必须串行化; - push 前必须检查 gRPC write queue readiness;
- 当 write queue not ready 时,连接记录 push queue block 时间戳,并以connection-busy语义失败;
- request/async push 分配 push request id,并在
RpcAckCallbackSynchronizer中注册 ack future; - bidirectional-stream 上收到的
Responsepayload 视为 push ack,用于清理或完成匹配的 future; - disconnect cleanup 必须清理该 connection 的 ack 状态。
需要强调:Push 和 ack 管道不定义领域订阅语义。Config、Naming 和 AI 规范定义应该推送什么以及何时推送,远程连接层只负责可靠地送达与确认。
7. 注销与 Listener 回调
当 transport terminated、active detection 失败、over-limit ejection 关闭连接,或其他服务端逻辑显式移除连接时,连接会被注销。规范定义的注销规则:
- listener 回调前必须先从注册表中移除连接;
- per-client-ip 计数必须递减,并在归零时移除;
- 底层 transport 必须关闭;
- 必须调用
ClientConnectionEventListener.clientDisConnected执行清理; - listener 失败应记录日志,但不得阻止其他 listener。
ConnectionManager.unregister()(ConnectionManager.java)严格遵循该顺序:先connections.remove(connectionId),再递减connectionForClientIp计数(归零时移除该 ip 条目),然后remove.close()关闭底层连接,最后notifyClientDisConnected通知 listener。
领域运行时 manager 可以通过clientConnected把状态挂载到connectionId,但必须通过clientDisConnected释放或重建这些状态。持久状态不得依赖连接生命周期——连接存在性只代表瞬时可达,不代表任何持久语义。
8. Active Detection 与运行时踢除
ConnectionManager启动周期性运行时踢除任务(@PostConstruct start(),见ConnectionManager.java第 255-279 行):启动 1 秒后每 3 秒执行一次runtimeConnectionEjector.doEject(),并同步更新长连接监控指标;模块维度连接计数监控由nacos.metric.grpc.server.connection.enabled(默认 true)与nacos.metric.grpc.server.connection.interval(默认 15 秒)控制。
默认执行器为NacosRuntimeConnectionEjector(NacosRuntimeConnectionEjector.java),其doEject()分两步:
public void doEject() { ejectOutdatedConnection(); // 过期连接检测 ejectOverLimitConnection(); // 过载踢除 }8.1 过期连接检测
- active time 超过
RuntimeConnectionEjector.KEEP_ALIVE_TIME的连接成为 active detection 候选; - push queue 在服务端 block 窗口内持续 block 的连接也成为候选(源码中窗口为
pushQueueBlockTimesLastOver(300 * 1000),即 300 秒); - active detection 发送
ClientDetectionRequest,等待成功响应(源码中单连接超时 5 秒,整体latch.await(5000L, TimeUnit.MILLISECONDS)); - 成功响应会刷新 active time(
connection.freshActiveTime()); - 未成功响应的连接会被
unregister注销。
8.2 过载踢除
- 只有设置了目标 load count 时才执行运行时 load ejection(
getLoadClient() > 0); - SDK 连接可以收到带可选 redirect address 的
ConnectResetRequest(ConnectionManager.loadSingle()会解析ip:port并设置 serverIp/serverPort,随后connection.request(connectResetRequest, 3000L)); - cluster 连接不得被 SDK load balancing 逻辑踢除(源码中
loadSingle只处理isSdkSource()的连接,ejectOverLimitConnection也只对 SDK 来源计数); - runtime ejector 实现可以通过 SPI 加载(
NacosServiceLoader.load(RuntimeConnectionEjector.class),名称来自 Control 配置ControlConfigs.getConnectionRuntimeEjector()),找不到匹配实现时回退到NacosRuntimeConnectionEjector,但必须保持ConnectionManager注册表和 listener 语义。
连接 control 规则与 TPS 检查是 Control 插件规范 定义的横切保护点。它们可以拒绝或延迟连接路径,但不得改变领域数据归属。
9. 领域使用方规则
使用远程连接生命周期的领域(Config、Naming、AI 等)必须遵循:
- 只把运行时状态绑定到
connectionId; - disconnect 时清理 connection-bound 状态;
- disconnect 处理必须幂等;
- 除非领域规范定义了可恢复身份,否则 reconnect 应视为新连接;
- 使用
RequestMeta中的 labels、source、能力表和 client version 做兼容决策,而不是重新解析 transport 内部细节; - 不得把连接存在性当作持久鉴权、归属或持久化证据;
- listener 回调中涉及慢速远端 IO 时,应避免阻塞连接注册或清理路径;
- push failure 和 connection-busy 行为必须通过领域 retry 或 resync 语义体现。
10. 相关规范
- 基础能力规范
- 集群成员规范
- 内部 RPC 与集群请求规范
- 请求过滤与运行时上下文规范
- gRPC API 规范
- Control 插件规范
- Trace 插件规范
小结
Nacos 的远程连接层把 gRPC 传输细节封装为一套统一、可观测、可保护的连接生命周期:AddressTransportFilter生成连接身份,ConnectionManager统一注册与注销,GrpcRequestAcceptor完成请求路由与上下文装配,RpcPushService与RpcAckCallbackSynchronizer保证推送可靠送达,NacosRuntimeConnectionEjector兜底清理异常连接。领域模块只需监听ClientConnectionEventListener事件并遵守绑定规则,即可在 Config、Naming、AI 等场景安全复用这套基础设施——这也是 Nacos 作为云原生动态服务发现与配置管理平台在连接治理层面的核心设计。
【免费下载链接】nacosan easy-to-use dynamic service discovery, configuration and service management platform for building AI cloud native applications.项目地址: https://gitcode.com/GitHub_Trending/na/nacos
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考