深入 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 后,理论上可以用 gRPC 客户端调用 Hadoop 服务(同样的 protobuf schema),但 SASL 认证、连接管理、错误处理仍然是 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 连接持续复用,多个 RPC 请求在连接上串行(连接级串行,不是请求级并行)。
每个 Client 实例维护一个 connection pool,按 (serverAddress, user) 维度缓存连接。如果同一个 Client 用同一个 user 调用同一个 server 的多个方法,复用同一条连接。
连接的生命周期:
1 | |
为什么连接级串行而不是并行?因为 Hadoop RPC 的 framing 是请求-响应模式——Client 发送一个请求,阻塞等待服务器响应,再发下一个。这种模式简化了实现(不需要处理乱序响应),代价是单连接的吞吐受限。
实际部署中,每个 Client JVM 可能并发向同一服务器发多个 RPC(多线程)。Client 用连接池处理这种并发——同一 (server, user) 的并发请求分到多条连接(最多 ipc.client.connection.size.pool.size 默认 10 条)。
服务端(NameNode / ResourceManager)接收大量客户端连接。NameNode 的 ipc server 默认 handler 线程数(dfs.namenode.handler.count)是 10,生产集群通常调到 100+。每个 handler 线程处理一条连接的请求队列。
Idle 清理机制
Hadoop RPC 客户端会主动清理空闲连接,避免连接泄漏。
ipc.client.connection.idlethreshold(默认 4000)控制连接空闲多久后被清理。一个 Client 调用 server 后如果 4 秒内没有新调用,连接被标记为 idle,下次 GC 时关闭。
为什么是 4 秒而不是更长?因为 Hadoop 的工作负载里,心跳型 RPC(DataNode → NameNode 心跳,默认 3 秒间隔)必须保持长连接。4 秒阈值让心跳连接不会被误清理(心跳间隔 3 秒 < 4 秒阈值)。
而长间隔的 RPC(例如 Client 偶尔调用 NameNode 的 getFileInfo)会建立新连接、用完就清理。这种"短连接复用 + idle 清理"模式让 NameNode 的连接表不会无限增长。
服务端也有连接清理。ipc.client.connection.idlethreshold 服务端版本控制 server 端 socket 的 idle 清理(默认与客户端一致 4 秒,但生产集群通常调大避免误清理)。
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
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.idlethreshold=10000 -ls / 这类命令行配置观察。修改 idle threshold 后连接保持时间变化,间接反映在 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 流 | 单连接串行 | 单连接串行 | 单连接串行 |
| 跨语言 | 强 | 强 | 弱(仅 Java/Python) | 中 |
| 版本兼容 | Protobuf optional | Thrift optional | 自定义 | 自定义 |
注意 gRPC 这一列。gRPC 用 HTTP/2 的多路复用让单连接可以并发多个请求,Hadoop RPC 是单连接串行。这是 Hadoop RPC 比 gRPC 吞吐低的一个原因(单连接并发度受限)。代价是 Hadoop RPC 实现更简单(不需要 HTTP/2 状态机)。
常见误解
误解一:“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 的心跳场景吞吐会怎么变化?需要做什么改造?
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 00-06 | HDFS 存储层 | 第一阶段(已完成) |
| 07-12 | YARN 资源管理层 | 第二阶段(已完成) |
| 13-15 | MapReduce 计算模型 | 第三阶段(已完成) |
| 16 | Hadoop RPC 协议栈:Protobuf、SASL 与 Connection Idle | 本篇 |
| 17 | 序列化与压缩:Writable、Avro 与 Codec | 下一篇 |
| 18 | Hadoop 安全:Kerberos、Delegation Token、Proxy User | |
| 19 | 监控与运维:Metrics V2、JMX、日志聚合 | |
| 20-22 | 演进、生态与对比 | 第五阶段,待开始 |
参考资料
- Apache Hadoop 官方文档:RPC. https://hadoop.apache.org/docs/current/hadoop-project-dist/hadoop-common/RpcEngine.html
- 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/
