深入 Hadoop 16 - Hadoop RPC 协议栈
前三阶段把 HDFS / YARN / MapReduce 三个子系统讲完了。第四阶段展开三个子系统共用的那套基础设施——RPC、序列化、安全、监控。本篇先讲 RPC。
Hadoop RPC 常被介绍成"Hadoop 进程间通信的协议"。这个描述对应了功能但完全没解释机制。准确的说法是:Hadoop RPC 是一个独立的、自研的 RPC 框架,不是 gRPC / Thrift / Dubbo 之类通用框架。它用 Protobuf 做序列化(Hadoop 2.x 起默认),用 SASL 做认证(Kerberos / Delegation Token / Simple),用 NIO + 长连接做传输。整个 Hadoop 项目(HDFS / YARN / MapReduce)的所有进程间通信都走这套 RPC——NameNode 与 DataNode、ResourceManager 与 NodeManager、ApplicationMaster 与 RM、Task 与 MRAppMaster。
本篇只抓一个问题:Hadoop RPC 的六层栈是怎么组织的、Protobuf 序列化怎么解决跨版本兼容、SASL 认证的三种模式各自适合什么场景、长连接的 idle 清理机制为什么这样设计。
为什么 Hadoop 不用通用 RPC 框架
Hadoop 项目 2006 年起步时,gRPC 还没出现(2015 年开源),Thrift 已经存在(2007 年 Facebook 开源)但 Hadoop 团队选择了自研。这个选择有三个动机:
第一,性能。Hadoop 集群里 NameNode 与数千 DataNode 的心跳、ResourceManager 与数千 NodeManager 的心跳、每个作业 AM 与 RM 的频繁 RPC,对 RPC 框架的吞吐和延迟要求极高。2006 年的 Thrift 性能不够好,自研的 Hadoop RPC 可以针对性优化(NIO 多路复用 + 长连接 + 简单 framing)。
第二,可控的版本兼容。Hadoop 集群升级通常是滚动升级——某段时间内集群里同时存在新旧版本的 NameNode / DataNode。RPC 协议必须向前向后兼容,让旧版本客户端能调用新版本服务器(通过 protobuf 的字段可选性)。
第三,与 Hadoop 安全深度集成。Hadoop 的 Kerberos 认证、Delegation Token 机制需要 RPC 层暴露认证回调。自研 RPC 可以让安全机制与传输层深度耦合,第三方框架做不到这种程度。
代价是 Hadoop RPC 是项目专有协议,不能与非 Hadoop 系统互通。Hadoop 2.x 起虽然使用 Protobuf 序列化,但传输 framing、RPC engine、SASL 协商和错误处理都不是 gRPC 协议,gRPC 客户端不能直接调用 Hadoop 服务。
六层 RPC 栈
Hadoop RPC 的内部结构可以拆成六层:
1 | |
每层都是可插拔的——Hadoop 配置文件可以指定使用哪种 RPC 引擎、哪种序列化、哪种 SASL 机制。这种可插拔设计让 Hadoop 可以在不同部署场景(高安全 / 高吞吐 / 简单)切换组合。
Protobuf 序列化与跨版本兼容
Hadoop 1.x 时代 RPC 序列化用 Writable——Hadoop 自己的 Java 序列化机制。每个 RPC 参数类必须实现 Writable 接口的 write(DataOutput) 和 readFields(DataInput) 方法。
Writable 的缺陷是版本兼容性差。如果某个 RPC 接口添加了一个新字段,旧版本客户端反序列化新版本服务器返回的对象时会失败(字段对不上)。Hadoop 集群滚动升级时这个问题很严重。
Hadoop 2.0 起 RPC 协议改用 Protobuf 定义。所有 RPC 接口在 .proto 文件里声明:
1 | |
Protobuf 的字段是 optional 的——新版本添加字段不会破坏旧版本反序列化(旧版本忽略未知字段)。这让 Hadoop 可以平滑滚动升级。
Hadoop 启动 ProtobufRpcEngine:
1 | |
实际上从 Hadoop 2.x 起,所有内部 RPC 接口(ClientNamenodeProtocol、DatanodeProtocol、ResourceTracker、ApplicationMasterProtocol 等)都强制使用 Protobuf。WritableRpcEngine 保留只是为了向后兼容旧 API(已不推荐使用)。
Protobuf 序列化的代价是性能——比 Writable 慢约 20%。但这个代价换来跨版本兼容,是值得的。Hadoop 3.x 起 JVM 优化(JIT 内联 Protobuf 的反序列化方法)让性能差距缩小到 10% 以内。
长连接与连接池
Hadoop RPC 客户端与服务器之间是长连接。一条 TCP 连接可同时承载多个 active call,每个请求由 callId 标识,响应可以乱序返回;同步调用只是调用线程等待结果,不代表 socket 上只能串行一个请求。
每个 Client 实例维护一个 connection pool,按 (serverAddress, user) 维度缓存连接。如果同一个 Client 用同一个 user 调用同一个 server 的多个方法,复用同一条连接。
连接的生命周期:
1 | |
Client 的 Connection 线程集中读取响应,再按 callId 唤醒对应调用者。并发线程可以共享同一 socket;异步 RPC 也复用这套 callId 与 active calls 表。
实际部署中,每个 Client JVM 可能并发向同一服务器发多个 RPC。连接池大小可以配置,但单条连接本身已经支持多路复用,不能把并发能力等同于连接条数。
服务端(NameNode / ResourceManager)接收大量客户端连接。NameNode 的 ipc server 默认 handler 线程数(dfs.namenode.handler.count)是 10,生产集群通常调到 100+。每个 handler 线程处理一条连接的请求队列。
Idle 清理机制
Hadoop RPC 客户端会主动清理空闲连接,避免连接泄漏。
ipc.client.idlethreshold(默认 4000)不是时间,而是连接数量阈值:连接数超过该值后,Client 才会检查并回收空闲连接。真正的连接空闲时长由 ipc.client.connection.maxidletime 控制,默认 10000 ms。
因此 4000 不能拿来与 3 秒心跳间隔比较,也不是服务端 socket 的 4 秒清理配置。偶发调用是否回收取决于连接数、maxidletime 和扫描时机。
SASL 认证的三种模式
Hadoop RPC 的 SASL 层支持三种认证模式:
1 | |
实际部署里这三种模式经常组合使用:
Client 第一次提交作业到 ResourceManager 时用 Kerberos 认证。RM 在 Kerberos 认证后签发一个 Delegation Token 给 Client,后续 Client 调用 RM 都用 Token(不需要再 Kerberos)。
作业启动后 AM 也拿到 RM 签发的 Token,AM 调用 RM(申请 container)和 HDFS(读写数据)都用 Token。
Token 有过期时间(默认 24 小时,可配置 dfs.namenode.delegation.token.renew-interval),长作业需要周期性 renew。
这种"Kerberos 换 Token"的设计平衡了安全与性能——Kerberos 保证初始身份验证的强安全性,Token 让长作业不需要频繁访问 KDC。
RPC 的版本协商
Hadoop RPC 客户端与服务器版本可能不同(滚动升级期间)。Protobuf 解决了消息格式的兼容性,但 RPC 协议本身还有版本协商:
1 | |
这种版本协商让 Hadoop 可以滚动升级——新版本 server 加入集群时旧版本 client 仍然能用。代价是 server 必须维护所有历史版本的协议代码(兼容代码累积,Hadoop 源码越来越庞大)。
实验:观察 Hadoop RPC
实验状态:UNVERIFIED_RUNTIME。下面是验证步骤,本轮没有连接真实 Hadoop 集群执行。
Hadoop RPC 暴露的运行时指标通过 JMX 访问:
1 | |
输出大量 JMX MBean,其中 RpcMetrics 和 RpcDetailedMetrics 暴露 RPC 层指标:
1 | |
这些指标对调试 RPC 性能很关键。CallQueueLength 持续高说明 server handler 不够(需要调大 handler 线程数)。RpcQueueTimeAvgTime 高说明 client 排队等 server 处理,瓶颈在 server CPU。
更细粒度的 RPC 调用统计在 RpcActivityForPort* MBean:
1 | |
客户端 RPC 行为可以通过 hadoop fs -D ipc.client.connection.maxidletime=30000 -ls / 这类命令行配置观察。修改最大空闲时间后,连接回收节奏会间接反映在 numOpenConnections 指标上。
抓包看 Hadoop RPC 协议:用 tcpdump 抓 NameNode 8020 端口流量,会看到大量小包(请求-响应配对)。每个 RPC 包格式是 4 字节长度 + Protobuf payload。Kerberos 模式下连接建立前还有 SASL 协商包(GSSAPI token exchange)。
模式提炼
Hadoop RPC 体现的设计模式:
1 | |
这个模式不只是 Hadoop。Cassandra 的内部 RPC(Native Transport)也是自研,基于 Netty + 自定义协议。Kafka 的协议是自研的 binary protocol,不依赖通用框架。这些系统都选择自研 RPC 因为通用框架(gRPC / Thrift)在性能、可控性、生态集成上不能完全满足。
工程迁移表
| Hadoop RPC 概念 | gRPC | Thrift | Cassandra Native | Kafka Protocol |
|---|---|---|---|---|
| 协议定义 | .proto | .thrift | 自定义 | 自定义 |
| 序列化 | Protobuf | TBinary / TCompact | 自定义 | 自定义二进制 |
| 传输 | HTTP/2 | TCP | TCP + Netty | TCP |
| 认证 | TLS / Token | SASL / TLS | SASL | SASL / TLS |
| 多路复用 | HTTP/2 流 | 依具体传输实现 | stream id | correlation id |
| 跨语言 | 强 | 强 | 弱(仅 Java/Python) | 中 |
| 版本兼容 | Protobuf optional | Thrift optional | 自定义 | 自定义 |
注意 gRPC 这一列。gRPC 用 HTTP/2 stream 做多路复用;Hadoop RPC 则用 callId 在自有 TCP framing 上复用并发调用。两者不互通,但差异不是“gRPC 并行、Hadoop 单连接串行”。
常见误解
误解一:“Hadoop RPC 用的是 gRPC”。Hadoop RPC 是自研的,不用 gRPC。即使 Hadoop 2.x 起用 Protobuf 做序列化,也是用 ProtobufRpcEngine 自己实现 RPC 引擎,不依赖 gRPC 库。
误解二:“Hadoop RPC 是无状态的”。RPC 协议本身是无状态的(每个 RPC 独立),但底层 TCP 连接是有状态的(长连接复用)。这种"协议无状态、传输有状态"是 RPC 系统的标准设计。
误解三:“Kerberos 认证每次 RPC 都查 KDC”。Kerberos 票据有生命周期(默认 10 小时),客户端拿到票据后缓存,期间所有 RPC 复用同一票据。只有票据过期才需要重新向 KDC 申请。Delegation Token 进一步减少 KDC 压力——Token 由 Hadoop 自己签发,不需要 KDC 参与。
误解四:“Hadoop RPC 是同步阻塞的”。客户端 API 看起来是同步阻塞的(调用立即返回结果),但底层用 NIO + Future 实现非阻塞 IO。多个 Client 线程可以共用一个连接的 NIO selector。
误解五:“Protobuf 序列化让 Hadoop RPC 慢”。Protobuf 比 Writable 慢约 10-20%,但绝对延迟仍然在毫秒级。RPC 性能瓶颈通常是 server handler 线程数和 callQueue 长度,不是序列化本身。把 ipc.server.handler 调大对吞吐的提升远大于换序列化方式。
练习
-
在 Hadoop 集群上访问
http://namenode-host:9870/jmx,找 RpcMetrics 和 RpcActivityForPort MBean。观察 NumOpenConnections、CallQueueLength、RpcQueueTimeAvgTime 三个指标随时间的变化(提交一个 MapReduce 作业前后对比)。 -
在
apache/hadoop源码里找到Client.java、Server.java、ProtobufRpcEngine.java、SaslRpcClient.java(hadoop-common-project/hadoop-ipc 模块),观察客户端连接池、服务端 handler 线程、Protobuf 序列化、SASL 认证的核心实现。 -
用 tcpdump 抓 NameNode 8020 端口的 TCP 流量,提交一个
hadoop fs -ls /,观察 Kerberos 协商(如果开启)+ ConnectionHeader + Protobuf 请求包的结构。 -
思考题:如果让 Hadoop RPC 切换到 gRPC(用 HTTP/2 多路复用),NameNode 与 1000 DataNode 的心跳场景吞吐会怎么变化?需要做什么改造?
系列导航
参考资料
- Apache Hadoop 源码:RPC.java. https://github.com/apache/hadoop/blob/rel/release-3.4.1/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RPC.java
- J. K. Ousterhout. Scheduling, Caching, and Distributed Systems.(关于 RPC 与连接管理的早期经典)
- Sanjay Radia. HADOOP-8990: Wire Protocol Compatibility.(Apache JIRA,Hadoop 2.x RPC 切换到 Protobuf 的设计文档)
- Apache Hadoop 源码:
Client.java、Server.java、ProtobufRpcEngine.java、SaslRpcClient.java. https://github.com/apache/hadoop - Konstantin Shvachko, Hairong Kuang. Hadoop RPC Performance.(Hadoop Summit 演讲,描述 Hadoop RPC 的性能模型)
- Google Protocol Buffers 文档:https://protobuf.dev/overview/
