前三阶段把 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
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 默认,使用 Java 原生序列化)
- 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.connection.idlethreshold 默认 4000)

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 连接持续复用,多个 RPC 请求在连接上串行(连接级串行,不是请求级并行)。

每个 Client 实例维护一个 connection pool,按 (serverAddress, user) 维度缓存连接。如果同一个 Client 用同一个 user 调用同一个 server 的多个方法,复用同一条连接。

连接的生命周期:

1
2
3
4
5
6
7
1. Client 第一次调用 server → 建立连接
- TCP 三次握手
- SASL 认证(Kerberos / Token)
2. 后续 RPC 在同一条连接上串行发送
3. 连接空闲超过 idle threshold → 被清理
- 默认 ipc.client.connection.maxidletime = 10000 ms (10 秒)
4. 连接断开(网络故障 / 服务器重启)→ 下次调用重建

为什么连接级串行而不是并行?因为 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
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 票据向服务器认证
- 强安全,但 KDC 压力大(每次 RPC 都要 TGS exchange)
- 第十八篇展开

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

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

1
2
# NameNode JMX(默认 8485 端口需要 Kerberos,简化用 50070 Web UI 的 jmx endpoint)
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.idlethreshold=10000 -ls / 这类命令行配置观察。修改 idle threshold 后连接保持时间变化,间接反映在 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 流 单连接串行 单连接串行 单连接串行
跨语言 弱(仅 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 调大对吞吐的提升远大于换序列化方式。

练习

  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-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 演进、生态与对比 第五阶段,待开始

参考资料