上一篇把 HDFS 切成 NameNode、DataNode、Client 三层职责。本篇展开 Client 写文件时的具体路径。

HDFS 写文件常被描述成"客户端把数据写到 HDFS"。这个描述把数据传输机制掩盖了。准确的说法是:客户端通过一个由 NameNode 提前挑选的 DataNode 链(Pipeline)以流水线方式逐包写入数据,每个包从链头流到链尾再返回 ACK,写完一个 block 后客户端向 NameNode 申请新 block。这种 Pipeline + ACK 队列的设计是 HDFS 写入路径的核心抽象,决定了 HDFS 写入吞吐、副本一致性和写入失败恢复的全部行为。

本篇只抓一个问题:客户端写到 HDFS 的字节是怎么一步步落到三个 DataNode 磁盘上的,以及为什么 HDFS 选择 Pipeline 而不是客户端并行写三份。

写入路径的六个步骤

把一次完整的 HDFS 写入展开,可以拆成六个步骤:

1
2
3
4
5
6
1. create() 客户端调用 NameNode 创建文件元数据
2. addBlock() 客户端向 NameNode 申请一个 block ID 和目标 DataNode 列表
3. Pipeline 建立:客户端连 DN-1,DN-1 连 DN-2,DN-2 连 DN-3
4. 数据分批发送:客户端把数据切成 64KB 的 packet,逐包发给 DN-1
5. ACK 回传:DN-3 收到后回 ACK 给 DN-2,DN-2 回给 DN-1,DN-1 回给客户端
6. close() 文件关闭时客户端通知 NameNode,NameNode 持久化文件

步骤 2 的细节最重要。NameNode 收到 addBlock 请求时,不是随便挑三个 DataNode 返回,而是按"副本放置策略"挑出最优的三个 DataNode。这个策略决定了写性能和数据可靠性,下文专门展开。

步骤 3 的 Pipeline 建立是 HDFS 写入路径最有特色的部分。客户端拿到三个 DataNode 列表后,不是分别连三个 DataNode 并行写数据,而是只连第一个 DataNode,让第一个 DataNode 连第二个 DataNode,第二个连第三个。数据形成一条流水线。

步骤 4 和 5 的 packet + ACK 队列机制让流水线既能压满带宽又不丢包。客户端把字节流切成 64KB 的 packet(默认 io.file.buffer.size 配置控制),每发一个 packet 就放进"ack queue"等待 ACK。ACK 返回后从 ack queue 移除,下一批发送。任何 packet 的 ACK 失败会触发 Pipeline 重构。

副本放置策略:为什么三个副本放两个机架

HDFS 默认副本数 3。NameNode 在 addBlock 时挑选三个 DataNode 的策略遵循一条经典规则(Hadoop 默认 BlockPlacementPolicyDefault):

1
2
3
4
第一个副本:放在写入客户端所在的 DataNode("本地写")
第二个副本:放在与第一个副本不同机架的某个 DataNode
第三个副本:放在与第二个副本相同机架、不同节点的某个 DataNode
后续副本(如果配置了副本数 > 3):随机放置

这条规则初看奇怪。三个副本只放在两个机架,第三个副本不放到第三个机架吗?

答案是写入吞吐和数据可靠性的权衡。如果把三个副本放在三个不同机架,写入时第二个和第三个副本都要跨机架(机架间带宽通常比机架内低),写入延迟变高。放在两个机架(本地机架一个、远端机架两个):

  • 第一个副本本地写(机架内,零跨机架带宽)
  • 第二个副本跨机架写一次(一次跨机架带宽)
  • 第三个副本在远端机架内写(机架内,零跨机架带宽)

整体只消耗一次跨机架带宽。同时可靠性仍然够:任意一个机架挂掉,至少还有一个机架持有两个副本(数据不丢)。任意两个机架同时挂掉的概率极低,可以接受。

机架感知策略依赖集群管理员在 core-site.xml 配置 topology.node.switch.mapping.script,提供一个脚本把 DataNode IP 映射到机架路径(例如 /rack-A/rack-A1)。NameNode 用这个映射计算"距离",决定副本放置。如果没有配置脚本,所有 DataNode 都被当成在同一个默认机架(/default-rack),副本放置策略会退化为随机选三个节点——这种部署在数据可靠性上是有损的。

副本放置策略是 HDFS 设计里少有的需要管理员配置的部分。生产集群必须配置机架脚本,否则任意一个机架故障可能丢数据。

Pipeline 写入:流水线的具体机制

客户端拿到 DataNode 列表后开始建立 Pipeline。以默认三副本为例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
Client                 DN-1                  DN-2                  DN-3
│ │ │ │
│ 1. TCP connect │ │ │
├────────────────────►│ │ │
│ │ 2. TCP connect │ │
│ ├────────────────────►│ │
│ │ │ 3. TCP connect │
│ │ ├────────────────────►│
│ │ │ │
│ 4. data packet p1 │ │ │
├────────────────────►│ 5. forward p1 │ │
│ ├────────────────────►│ 6. forward p1 │
│ │ ├────────────────────►│
│ │ │ │
│ │ │ 7. write to disk │
│ │ 8. ack p1 │ │
│ │◄────────────────────┤ │
│ 9. ack p1 │ │ │
│◄────────────────────┤ │ │

注意步骤 4-9 之间,packet 是流水线方式转发的——DN-1 收到 p1 之后立刻把 p1 转发给 DN-2,自己继续接收客户端的 p2。整条链路在稳定状态下同时承载三个 packet(一个在 Client → DN-1 段、一个在 DN-1 → DN-2 段、一个在 DN-2 → DN-3 段),这就是 Pipeline 的吞吐优势。

为什么 HDFS 选 Pipeline 而不是客户端并行写三份?两个原因:

第一,客户端出网带宽。客户端如果把数据并行写到三个 DataNode,客户端出网带宽消耗三倍。Pipeline 模式下客户端只发一次,DN-1 转发给 DN-2,DN-2 转发给 DN-3——每个 DataNode 只承担一次出网带宽。在千兆网卡和 1000 节点集群里,这个差异决定客户端能不能压满网卡。

第二,副本一致性。Pipeline 模式下 packet 顺序到达 DN-3 的顺序就是客户端发送顺序——DN-1 转发是 FIFO,DN-2 转发也是 FIFO。如果客户端并行写三份,三个 DataNode 上 packet 到达顺序可能不一致,需要额外的序号机制保证一致。Pipeline 把序号约束变成 TCP 字节流的天然属性,简化了一致性管理。

数据队列与 ACK 队列

客户端写数据时内部维护两个队列:

1
2
data queue:待发送的 packet 队列(按生成顺序)
ack queue:已发送、等待 ACK 的 packet 队列(按发送顺序)

DFSOutputStream 的标准工作循环:

  1. 应用层调用 write(byte[]),把字节攒到内部 buffer
  2. buffer 攒够一个 packet 大小(默认 64KB)后,DFSOutputStream 把 packet 入 data queue 尾部
  3. 一个 Sender 线程从 data queue 头部取 packet,通过 Pipeline 发给 DN-1,packet 出 data queue 入 ack queue
  4. 一个 ResponseProcessor 线程等待 Pipeline 返回的 ACK
  5. 收到 ACK 后,ResponseProcessor 把对应 packet 从 ack queue 移除
  6. 任意 packet 的 ACK 失败时,整个 Pipeline 进入恢复流程(下文展开)

这种"两个队列 + 两个线程"的设计让客户端既能压满 Pipeline 带宽(Sender 持续发),又能保证可靠性(ResponseProcessor 持续确认)。任何 packet 失败时,ack queue 里所有未 ACK 的 packet 都要重发。

注意一个细节:客户端的 write() 调用返回不代表数据已经持久化。在 ack queue 里的 packet 没收到 ACK 之前,对应字节都不算"已写入 HDFS"。如果客户端调用了 write() 但 Pipeline 出现故障导致 ack queue 里的 packet 没有成功写入,应用层会看到 IOException。这是 HDFS 写操作的可靠性语义。

Pipeline 故障恢复

Pipeline 写入过程中任意一个 DataNode 故障(磁盘满、网络断、进程崩),Pipeline 必须恢复。HDFS 的恢复流程是这套写入路径设计里最精巧的部分。

故障检测:客户端的 ResponseProcessor 在等待 ACK 时超时,或者 Pipeline socket 抛 IOException。客户端把 ack queue 里所有 packet 标记为待重发。

故障定位:客户端通过 ACK 协议判断哪个 DataNode 故障。每个 DataNode 在转发 packet 时会附带自己的序号,如果某个 DataNode 序号停滞,客户端就知道是它出问题。客户端把这个 DataNode 从 Pipeline 里移除。

Pipeline 重构:客户端用剩下的健康 DataNode 重新建立 Pipeline。例如三副本 Pipeline 中 DN-2 故障,新 Pipeline 是 Client → DN-1 → DN-3。

同步块状态:客户端向 NameNode 报告"DN-2 在写某个 block 时出问题",NameNode 在内存里把这个 block 的副本数从 3 标记为 2,安排后台任务在其他 DataNode 上补齐第三个副本。同时客户端继续往新 Pipeline 写剩余数据,block ID 不变。

Block 完成与校验:block 写完后,客户端告诉 NameNode 这个 block 已经关闭。NameNode 等到该 block 的 blockReport 显示三个副本都到位(或者足够时间后)才认为 block 完全持久化。

这种"故障时丢弃坏节点、用健康节点继续写、后台补齐副本"的恢复流程,是 HDFS 假设"节点总会失败"前提下的工程对策。整个流程对客户端应用透明——客户端只需要捕获 IOException 重试,不需要知道 Pipeline 是怎么重构的。

实验:观察 Pipeline 写入

可以用 hdfs dfs -put 上传一个本地文件,同时观察 DataNode 的日志或 jstack 来看到 Pipeline 的具体行为。

准备一个本地大文件(100MB 左右,足以分成多个 128MB block 的话用一个就够,但跨多个 block 更有观察价值):

1
2
dd if=/dev/urandom of=/tmp/test-200mb.bin bs=1M count=200
hdfs dfs -put /tmp/test-200mb.bin /test/

写入时打开 NameNode Web UI 的 Startup Progress 或者 Datanodes 标签页,可以看到写入过程中某些 DataNode 的 Block Pool Used 数字增长。这正是 Pipeline 数据落盘的证据。

更直接的观察方式是用 hdfs dfs -fsck /test/test-200mb.bin -files -blocks -locations

1
2
3
4
5
6
7
8
9
/test/test-200mb.bin 209715200 bytes, 2 block(s), writing...
0. BP-1234567890-namenode-host:blk_1073741825_1001 \
len=134217728 Live_repl=3 [DatanodeInfoWithStorage[10.0.0.1:9866,DS-...], \
DatanodeInfoWithStorage[10.0.0.2:9866,DS-...], \
DatanodeInfoWithStorage[10.0.0.3:9866,DS-...]]
1. BP-1234567890-namenode-host:blk_1073741826_1002 \
len=75497472 Live_repl=3 [DatanodeInfoWithStorage[10.0.0.2:9866,DS-...], \
DatanodeInfoWithStorage[10.0.0.3:9866,DS-...], \
DatanodeInfoWithStorage[10.0.0.4:9866,DS-...]]

这个输出直接展示了三个事实:

文件 200MB 被切成 2 个 block,第一个 128MB、第二个 72MB(剩余字节)。

每个 block 有 3 个副本(Live_repl=3),分布在不同的 DataNode 上。

两个 block 的副本放置遵循"本地写 + 跨机架 + 跨机架同节点"规则。如果配置了机架脚本,可以看到第一个副本所在 DataNode 与上传客户端所在节点是同一个。

如果想让 Pipeline 故障恢复可见,可以故意杀掉 Pipeline 中间的 DataNode 进程,再上传文件。客户端日志里会出现 “Exception in createBlockOutputStream” + “Abandoning block” + “Connecting to datanode” 这类重试信息。NameNode Web UI 上对应 block 的 Under Replicated Blocks 会瞬间涨 1,几秒后系统调度新副本归零。

模式提炼

HDFS 写入路径体现的设计模式:

1
2
3
4
5
6
7
模式:流水线写入 + ACK 队列 + 故障节点剔除

- 把数据切成定长 packet,让 packet 沿一条 DataNode 链流水转发
- 客户端只发一份,每个 DataNode 转发给下一个,避免客户端出网带宽放大
- 每发一个 packet 入 ACK 队列,收到链尾 ACK 后出队
- 任意节点失败时把它从链里剔除,用剩下的健康节点重新建链
- 故障节点上未完成的 block 由后台任务在其他节点补齐副本

这个模式不只在 HDFS 出现。Ceph 的 RBD 写入也用 Pipeline-like 链式复制(主 OSD 收到后转发给从 OSD)。MySQL 半同步复制是 Pipeline 的退化版(主从两节点)。Kafka 的 ISR 写入方式(leader 收到、follower 拉取)虽然不是严格 Pipeline,但同样遵循"主节点统一接收、副本按序同步"的思路,差别在 Kafka 用 pull 模型、HDFS 用 push 模型。

工程迁移表

HDFS 写入概念 GFS Ceph Kafka MySQL 半同步
副本放置策略 Master 选 ChunkServer MON 计算 PG → OSD 映射 Controller 分配 Partition → Broker DBA 配置 replication
Pipeline 链式转发 是(PG 主 → 从) 否(leader 推 / follower 拉) 否(binlog 单向复制)
ACK 队列 Operation Log PG journal ISR commit log binlog position
packet 大小 64KB 4MB 默认 message set binlog event
故障节点剔除 是(PG peering) ISR 收缩 master/slave 切换
副本补齐 Master 重备份 MON peering ISR 扩展 重建 slave
顺序一致性 Pipeline FIFO PG journal 顺序 ISR offset binlog position

注意 Kafka 这一列的差异。Kafka 的写入不是严格的 Pipeline——生产者写 leader,follower 主动从 leader 拉取。这种"主写从拉"模式让 Kafka 可以容忍 follower 慢(拉不动就暂时落后,不影响 leader 写入性能),代价是 follower 落后时一致性窗口扩大。HDFS 的 Pipeline 是同步模式——三个副本必须全部 ACK 才算成功,慢节点会拖慢整体写入,但一致性窗口几乎为零。

常见误解

误解一:“HDFS 写入是同步等待三个副本落盘”。不完全是。Pipeline 在 packet 级别是同步的(每个 packet 必须三个副本 ACK 才算成功),但 packet 落盘到 DataNode 磁盘时是异步的——DataNode 收到 packet 后先入内存 buffer,定期刷盘。如果 DataNode 进程崩溃,buffer 里的 packet 可能丢。HDFS 通过 fsync 在 block 关闭时强制刷盘缓解这个问题,但写过程中仍存在数据丢失窗口。

误解二:“HDFS 写入失败客户端要手动重试”。部分正确。Pipeline 中某个 DataNode 失败时客户端会自动剔除该节点继续写,不需要应用层重试。但如果整个 Pipeline 都失败(例如客户端到 NameNode 都连不通),应用层会收到 IOException,需要重试。这种"局部失败自动恢复、整体失败抛异常"是 HDFS 的标准错误语义。

误解三:“副本数越高数据越安全”。副本数高确实降低数据丢失概率,但写入吞吐会成比例下降。三副本写一次的字节量是数据本身的 3 倍。如果副本数提到 5,写入吞吐大约降为 60%。生产环境通常用 3 副本配合机架感知,比单纯提高副本数更划算。Hadoop 3.x 引入纠删码(第六篇展开)是另一种思路——用计算换存储,副本开销从 3x 降到 1.5x。

误解四:“Pipeline 写入顺序就是文件字节顺序”。在单个 block 内是对的——Pipeline 保证 packet 顺序到达 DN-3 与客户端发送顺序一致。但跨 block 时不一定:客户端可以并发申请多个 block 并发 Pipeline 写入(特别在大文件场景),不同 block 的字节落盘顺序可能与应用层调用顺序不完全一致。HDFS 通过客户端层面的单 DFSOutputStream 串行化缓解了这个问题。

误解五:“NameNode 知道每个字节当前写到哪个 DataNode 了”。NameNode 不知道。NameNode 只在 addBlock 时返回 DataNode 列表,之后整个 Pipeline 的数据传输都不经过 NameNode。NameNode 知道 block 已经分配、知道 block 已经关闭、知道 block 当前在哪些 DataNode(通过 blockReport),但不知道 Pipeline 此刻传输的具体字节。

练习

  1. hdfs dfs -put 上传一个 300MB 的本地文件到 HDFS,再用 hdfs fsck /path/to/file -files -blocks -locations 查看 block 分布。验证每个 block 的三个副本是否符合"本地写 + 跨机架 + 跨机架同节点"放置规则(前提是集群配置了机架脚本)。

  2. apache/hadoop 源码里找到 DataStreamer.javaResponseProcessor.java(hadoop-hdfs-client 模块下),观察 data queue 和 ack queue 的具体实现。这两个内部类是 HDFS Pipeline 写入的核心。

  3. 在伪分布式集群上故意 kill 一个 DataNode 进程,立刻 hdfs dfs -put 上传文件,观察客户端日志里的 Pipeline 恢复过程。然后重启该 DataNode,用 hdfs dfsadmin -report 观察 Under Replicated Blocks 数字的变化(应该瞬间涨、几分钟后归零)。

  4. 思考题:如果把 Pipeline 模式改成客户端并行写三个 DataNode(每个副本独立写),HDFS 写入吞吐会变好还是变差?数据一致性会发生什么变化?

系列导航

序号 主题 状态
00 导读:Hadoop 的核心前提是节点总会失败
01 HDFS 架构:NameNode、DataNode 与元数据的三层切分 上一篇
02 文件写入路径:从客户端到 DataNode 的流水线 本篇
03 文件读取路径:副本选择与短路读 下一篇
04 NameNode 内存模型:FSImage、EditLog 与启动恢复
05 HDFS HA:Quorum Journal Manager、ZKFC 与脑裂防御
06 HDFS 3.x 演进:纠删码、路由联邦与 Observer NameNode

参考资料

  • Sanjay Ghemawat, Howard Gobioff, Shun-Tak Leung. The Google File System. SOSP 2003. Section 3.2 “System Interactions” 描述了 Pipeline 写入和 ack 队列机制。
  • Apache Hadoop 官方文档:HDFS Architecture. https://hadoop.apache.org/docs/current/hadoop-project-dist/hadoop-hdfs/HdfsDesign.html
  • Apache Hadoop 源码:DataStreamer.javaDFSOutputStream.java. https://github.com/apache/hadoop
  • Konstantin Shvachko. Hadoop Swap Space: Block Placement.(Hadoop Summit 演讲,详细解释了三副本放置策略的机架感知)
  • Tom White. Hadoop: The Definitive Guide. O’Reilly, 4th Edition 2015. Chapter 3 详细描述了 Pipeline 写入的具体步骤和 packet 流转。