前三阶段把 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
1. 应用层(Application Layer)
- 用户代码调用 ClientNamenodeProtocol.create() 等 RPC 接口
- 接口由 .proto 文件定义,Hadoop 编译时生成 Java 接口和桩

2. RPC 引擎层(RPC Engine)
- WritableRpcEngine(Hadoop 1.x 默认,使用 Hadoop Writable 协议)
- ProtobufRpcEngine(Hadoop 2.x 起默认,使用 Protobuf 序列化)
- 用户可自定义(如 AvroRpcEngine)
- 负责:序列化参数、派发到服务端实现、反序列化响应

3. 序列化层(Serialization)
- Writable(Hadoop 原生,要求实现 write/readFields 方法)
- Protobuf(Hadoop 2.x 起主流,跨语言跨版本)
- Avro(Hadoop 生态系统常用,例如 Hive / Pig 内部)
- 下一篇文章展开

4. SASL 认证层
- Kerberos(GSSAPI SASL 机制)
- Delegation Token(DIGEST-MD5 SASL 机制)
- Simple(无认证,仅靠 IP)
- 第十八篇展开

5. 传输层(Transport)
- NIO 多路复用 + 长连接(每对 client-server 一条 TCP)
- 长度前缀 framing(4 字节 length + payload)
- 客户端连接池(`ipc.client.idlethreshold` 默认 4000,表示连接数超过该阈值后才检查 idle)

6. TCP/IP
- 默认 keepalive 开启
- TCP_NODELAY 开启(禁用 Nagle,让小 RPC 立即发送)

每层都是可插拔的——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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
service ClientNamenodeProtocol {
rpc create(CreateRequestProto) returns (CreateResponseProto);
rpc append(AppendRequestProto) returns (AppendResponseProto);
// ... 几十个 RPC 方法
}

message CreateRequestProto {
required string src = 1;
required FsPermissionProto masked = 2;
required string clientName = 3;
required int32 flag = 4;
optional bool createParent = 5;
required int32 replication = 6;
required int64 blockSize = 7;
// ... 后续可以添加 optional 字段,不破坏旧客户端
}

Protobuf 的字段是 optional 的——新版本添加字段不会破坏旧版本反序列化(旧版本忽略未知字段)。这让 Hadoop 可以平滑滚动升级。

Hadoop 启动 ProtobufRpcEngine:

1
2
3
4
<property>
<name>hadoop.rpc.engine</name>
<value>org.apache.hadoop.ipc.ProtobufRpcEngine</value>
</property>

实际上从 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
2
3
4
5
6
7
1. Client 第一次调用 server → 建立连接
- TCP 三次握手
- SASL 认证(Kerberos / Token)
2. 后续 RPC 在同一条连接上按 callId 多路复用
3. 连接空闲超过 idle threshold → 被清理
- 默认 ipc.client.connection.maxidletime = 10000 ms (10 秒)
4. 连接断开(网络故障 / 服务器重启)→ 下次调用重建

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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
1. SIMPLE(hadoop.security.authentication=simple)
- 无认证,靠客户端 IP + 用户名(whoami 输出)
- 适合测试集群、单租户开发环境
- 生产集群不允许

2. KERBEROS(hadoop.security.authentication=kerberos)
- SASL GSSAPI 机制
- 客户端用 Kerberos 票据向服务器认证
- 强安全;服务票据可在 credential cache 中复用,并非每次 RPC 都访问 KDC
- 第十八篇展开

3. TOKEN(实际是 KERBEROS 模式的派生)
- SASL DIGEST-MD5 机制
- 客户端用 Delegation Token 认证(不是 Kerberos 票据)
- Token 由 NameNode / ResourceManager 在 Kerberos 认证后签发
- 减轻 KDC 压力(Kerberos 一次,Token 多次使用)
- 适合长作业(MapReduce / Spark 持续几小时)

实际部署里这三种模式经常组合使用:

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
2
3
4
5
6
7
8
9
10
11
12
Client 连接 Server 时发送 ConnectionHeader:
- 协议名(例如 "org.apache.hadoop.hdfs.protocol.ClientNamenodeProtocol")
- 协议版本(例如 1L)
- 认证信息(ugi, auth method)

Server 检查协议名和版本是否支持:
- 支持 → 接受连接
- 不支持 → 拒绝连接,返回错信息

如果 server 支持协议但版本高于 client(server 新):
- server 接受连接,按 client 版本回应(向后兼容)
- protobuf 的 optional 字段保证消息兼容

这种版本协商让 Hadoop 可以滚动升级——新版本 server 加入集群时旧版本 client 仍然能用。代价是 server 必须维护所有历史版本的协议代码(兼容代码累积,Hadoop 源码越来越庞大)。

实验:观察 Hadoop RPC

实验状态:UNVERIFIED_RUNTIME。下面是验证步骤,本轮没有连接真实 Hadoop 集群执行。

Hadoop RPC 暴露的运行时指标通过 JMX 访问:

1
2
# NameNode JMX(Hadoop 3.4.1 的 NameNode HTTP Web UI 默认端口是 9870;开启安全后按集群 SPNEGO/HTTPS 配置访问)
curl http://namenode-host:9870/jmx

输出大量 JMX MBean,其中 RpcMetricsRpcDetailedMetrics 暴露 RPC 层指标:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
{
"beans": [
{
"name": "Hadoop:service=NameNode,name=RpcMetrics",
"RpcQueueTimeNumOps": 1523456, // RPC 总数
"RpcQueueTimeAvgTime": 0.5, // 平均排队时间(ms)
"RpcProcessingTimeNumOps": 1523456,
"RpcProcessingTimeAvgTime": 2.3, // 平均处理时间(ms)
"NumOpenConnections": 234, // 当前活跃连接数
"CallQueueLength": 5, // 当前请求队列长度
"ReceivedBytes": 1234567890,
"SentBytes": 9876543210
}
]
}

这些指标对调试 RPC 性能很关键。CallQueueLength 持续高说明 server handler 不够(需要调大 handler 线程数)。RpcQueueTimeAvgTime 高说明 client 排队等 server 处理,瓶颈在 server CPU。

更细粒度的 RPC 调用统计在 RpcActivityForPort* MBean:

1
2
3
4
5
6
{
"name": "Hadoop:service=NameNode,name=RpcActivityForPort8020",
"callQueueLength": 3,
"numDroppedConnections": 0,
"numOpenConnections": 234
}

客户端 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
2
3
4
5
6
7
8
模式:自研轻量 RPC + 可插拔序列化 + 长连接 + SASL 认证

- 不依赖通用 RPC 框架(gRPC / Thrift),自研针对场景优化
- 序列化层可插拔(Writable / Protobuf / Avro),换序列化不需要改业务代码
- 用 SASL 抽象认证层,支持多种机制(Kerberos / Token / Simple)
- 客户端长连接复用,连接池 + idle 清理避免泄漏
- 服务端 handler 线程池处理并发请求
- Protobuf optional 字段保证滚动升级兼容性

这个模式不只是 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 调大对吞吐的提升远大于换序列化方式。

练习

  1. 在 Hadoop 集群上访问 http://namenode-host:9870/jmx,找 RpcMetrics 和 RpcActivityForPort MBean。观察 NumOpenConnections、CallQueueLength、RpcQueueTimeAvgTime 三个指标随时间的变化(提交一个 MapReduce 作业前后对比)。

  2. apache/hadoop 源码里找到 Client.javaServer.javaProtobufRpcEngine.javaSaslRpcClient.java(hadoop-common-project/hadoop-ipc 模块),观察客户端连接池、服务端 handler 线程、Protobuf 序列化、SASL 认证的核心实现。

  3. 用 tcpdump 抓 NameNode 8020 端口的 TCP 流量,提交一个 hadoop fs -ls /,观察 Kerberos 协商(如果开启)+ ConnectionHeader + Protobuf 请求包的结构。

  4. 思考题:如果让 Hadoop RPC 切换到 gRPC(用 HTTP/2 多路复用),NameNode 与 1000 DataNode 的心跳场景吞吐会怎么变化?需要做什么改造?

系列导航

序号 主题 状态
00 导读:节点总会失败
01 HDFS 架构与三层切分
02 文件写入路径与流水线
03 文件读取路径与副本选择
04 NameNode 内存模型与启动恢复
05 HDFS HA 与脑裂防御
06 HDFS 3.x 演进与纠删码
07 YARN 架构与三方契约
08 YARN 资源模型、Container 与 NodeLabel
09 YARN 调度器对比:FIFO、Capacity、Fair
10 YARN 应用程序生命周期
11 YARN HA 与 Federation
12 YARN Timeline Service v2
13 MapReduce 编程模型与分而治之
14 MapReduce Shuffle 全流程
15 MRv2 on YARN ApplicationMaster 与 Task Attempt 上一篇
16 Hadoop RPC 协议栈 本篇
17 序列化与压缩 下一篇
18 Hadoop 安全 Kerberos Token ProxyUser
19 监控与运维 Metrics JMX 日志聚合
20 Hadoop 生态 Hive HBase Pig
21 Hadoop 与对象存储 Kubernetes 演进对比
22 Hadoop 设计遗产从 GFS MapReduce 到云原生

参考资料