深入 Hadoop 02 - 文件写入路径与流水线
上一篇把 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 的细节最重要。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 | |
这条规则初看奇怪。三个副本只放在两个机架,第三个副本不放到第三个机架吗?
答案是写入吞吐和数据可靠性的权衡。如果把三个副本放在三个不同机架,写入时第二个和第三个副本都要跨机架(机架间带宽通常比机架内低),写入延迟变高。放在两个机架(本地机架一个、远端机架两个):
- 第一个副本本地写(机架内,零跨机架带宽)
- 第二个副本跨机架写一次(一次跨机架带宽)
- 第三个副本在远端机架内写(机架内,零跨机架带宽)
整体只消耗一次跨机架带宽。同时可靠性仍然够:任意一个机架挂掉,至少还有一个机架持有两个副本(数据不丢)。任意两个机架同时挂掉的概率极低,可以接受。
机架感知策略依赖集群管理员在 core-site.xml 配置 topology.node.switch.mapping.script,提供一个脚本把 DataNode IP 映射到机架路径(例如 /rack-A/rack-A1)。NameNode 用这个映射计算"距离",决定副本放置。如果没有配置脚本,所有 DataNode 都被当成在同一个默认机架(/default-rack),副本放置策略会退化为随机选三个节点——这种部署在数据可靠性上是有损的。
副本放置策略是 HDFS 设计里少有的需要管理员配置的部分。生产集群必须配置机架脚本,否则任意一个机架故障可能丢数据。
Pipeline 写入:流水线的具体机制
客户端拿到 DataNode 列表后开始建立 Pipeline。以默认三副本为例:
1 | |
注意步骤 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 | |
DFSOutputStream 的标准工作循环:
- 应用层调用
write(byte[]),把字节攒到内部 buffer - buffer 攒够一个 packet 大小(默认 64KB)后,DFSOutputStream 把 packet 入 data queue 尾部
- 一个 Sender 线程从 data queue 头部取 packet,通过 Pipeline 发给 DN-1,packet 出 data queue 入 ack queue
- 一个 ResponseProcessor 线程等待 Pipeline 返回的 ACK
- 收到 ACK 后,ResponseProcessor 把对应 packet 从 ack queue 移除
- 任意 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 | |
写入时打开 NameNode Web UI 的 Startup Progress 或者 Datanodes 标签页,可以看到写入过程中某些 DataNode 的 Block Pool Used 数字增长。这正是 Pipeline 数据落盘的证据。
更直接的观察方式是用 hdfs dfs -fsck /test/test-200mb.bin -files -blocks -locations:
1 | |
这个输出直接展示了三个事实:
文件 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 | |
这个模式不只在 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 此刻传输的具体字节。
练习
-
用
hdfs dfs -put上传一个 300MB 的本地文件到 HDFS,再用hdfs fsck /path/to/file -files -blocks -locations查看 block 分布。验证每个 block 的三个副本是否符合"本地写 + 跨机架 + 跨机架同节点"放置规则(前提是集群配置了机架脚本)。 -
在
apache/hadoop源码里找到DataStreamer.java和ResponseProcessor.java(hadoop-hdfs-client 模块下),观察 data queue 和 ack queue 的具体实现。这两个内部类是 HDFS Pipeline 写入的核心。 -
在伪分布式集群上故意 kill 一个 DataNode 进程,立刻
hdfs dfs -put上传文件,观察客户端日志里的 Pipeline 恢复过程。然后重启该 DataNode,用hdfs dfsadmin -report观察 Under Replicated Blocks 数字的变化(应该瞬间涨、几分钟后归零)。 -
思考题:如果把 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.java、DFSOutputStream.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 流转。
