HDFS架构设计与源码导读
发布日期: 2026/07/31 阅读总量: 0

一、真实场景:NameNode Full GC 30分钟

2019年我维护一个50节点Hadoop 2.7.3集群,某天业务高峰期NameNode突发Full GC,持续30分钟,期间客户端无法创建/读取文件,Hive作业全部失败。日志显示GC overhead limit exceeded,原因在于堆内存中元数据对象过多,且RPC请求堆积导致CPU飙高。

故障后分析:NameNode管理了约3亿个文件/目录对象(每个约150字节),堆内存32GB不够;RPC线程池处理请求时大量等待获取读锁,GC停顿后锁竞争更严重。这个问题直接暴露了HDFS中心化架构的两大弱点:单点瓶颈元数据膨胀

二、问题:中心化元数据管理为什么是瓶颈?

HDFS采用NameNode + DataNode主从架构。NameNode在内存中维护整个文件系统的目录树和文件块映射,所有增删改查都要经过它。这种设计好处是简单可靠,但存在天花板:

  • 内存容量限制:每个文件/目录/块信息占用约150~200字节,10亿个对象需要20GB+,受限于单机内存
  • RPC吞吐限制:NameNode单机处理能力约20~30K QPS,超过后请求排队
  • GC停顿敏感:老年代对象过多,Full GC次数增多,停顿时间秒级
  • 单点不可用:NameNode宕机则集群不可用(除非HA)

三、方案对比:HA vs Federation vs 纠删码

针对上述问题,Hadoop社区提出了三种进化方案。我基于Hadoop 3.2.1版本进行实际部署对比。

方案1:HDFS HA(Active/Standby)

通过两个NameNode(Active+Standby)共享元数据在JournalNode上的EditLog,实现秒级故障切换。解决了单点不可用,但Active节点仍为单点瓶颈。元数据仍只能放在一台NameNode内存中。

方案2:HDFS Federation

多个独立的NameNode各自管理部分目录(如/namespace1、/namespace2),每个NameNode有自己的块池。DataNode同时注册到所有NameNode。这样元数据水平扩展,内存压力分散,RPC也分散。但需要客户端指定集群的映射,运维复杂度增加。

方案3:Erasure Coding (纠删码)

将3副本改为1.5倍存储的纠删码,实质是存储优化,不解决元数据瓶颈,但能降低DataNode磁盘用量。可与Federation/HA组合。

对比表:

维度 标准架构 HA Federation 纠删码
高可用 秒级切换 依赖各NameNode自身HA 无影响
元数据容量扩展 只能单机 只能单机 水平扩展 无影响
RPC吞吐扩展 单点 单点 水平扩展 无影响
运维复杂度 中(JournalNode) 高(目录规划) 中(编码策略)
存储效率 3倍 3倍 3倍 1.1~1.5倍

四、完整代码实现:从客户端到核心源码

4.1 Java客户端写入HDFS

以下代码在Hadoop 3.2.1,JDK 8下测试通过,可写入100MB文件。

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataOutputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import java.io.IOException;

public class HdfsWriteClient {
    public static void main(String[] args) throws IOException {
        Configuration conf = new Configuration();
        conf.set("fs.defaultFS", "hdfs://namenode:8020");
        conf.set("dfs.client.use.datanode.hostname", "true");
        conf.set("dfs.replication", "3");
        FileSystem fs = FileSystem.get(conf);
        Path filePath = new Path("/tmp/test_100m.dat");
        // 写入100MB随机数据
        byte[] data = new byte[100 * 1024 * 1024];
        new java.util.Random().nextBytes(data);
        FSDataOutputStream out = fs.create(filePath, true);
        out.write(data);
        out.hflush();
        out.close();
        System.out.println("写入完成: " + filePath);
        fs.close();
    }
}

4.2 RPC协议定义(protobuf)

NameNode和客户端通过Hadoop IPC通信,协议定义在hadoop-common-project/hadoop-common/src/main/proto/ClientNamenodeProtocol.proto中(以2.8+版本为例)。展示创建文件请求的核心proto:

message CreateRequestProto {
    required string src = 1;        // 文件路径
    required FsPermissionProto masked = 2; // 权限
    required string clientName = 3;
    optional CreateFlagProto createFlag = 4 [default = CREATE];
    optional bool createParent = 5 [default = true];
    optional uint32 replication = 6 [default = 3];
    optional uint64 blockSize = 7 [default = 134217728];
    repeated CryptoProtocolVersionProto cryptoProtocolVersion = 8;
    optional string ecPolicyName = 9;
    optional string storagePolicy = 10;
    optional bool isLazyPersist = 11 [default = false];
}
message CreateResponseProto {
    optional HdfsFileStatusProto fs = 1; // 返回文件状态
}

4.3 NameNode RPC处理核心:FSNamesystem.createFile

客户端调用ClientProtocol.create(),实际由NameNodeRpcServer转发给FSNamesystem.createFile()。以下是简化后的关键逻辑(基于Hadoop 3.2.1源码):

// FSNamesystem.java
public HdfsFileStatus startFile(String src, PermissionStatus permissions,
                                String holder, String clientMachine,
                                EnumSet<CreateFlag> flag, boolean createParent,
                                short replication, long blockSize,
                                CryptoProtocolVersion[] cryptoProtocolVersions,
                                String ecPolicyName, String storagePolicy,
                                boolean isLazyPersist)
    throws IOException {
    // 1. 写锁
    writeLock();
    try {
        // 2. 检查路径是否存在
        FSDirectory dir = getFSDirectory();
        INodesInPath iip = dir.getINodesInPath(src, false);
        if (iip.isRoot()) throw new FileAlreadyExistsException(src);
        // 3. 分配inodeId
        long inodeId = dir.allocateNewInodeId();
        // 4. 创建INodeFile并插入目录树
        INodeFile newNode = new INodeFile(inodeId, ...);
        dir.addChild(iip, newNode);
        // 5. 记录EditLog
        getEditLog().logOpenFile(src, newNode);
        // 6. 返回文件状态(未分配任何块)
        return new HdfsFileStatus(...);
    } finally {
        writeUnlock();
    }
}

4.4 DataNode写入Pipeline

当客户端写入第一个块时,NameNode返回三个DataNode列表,客户端构建pipeline。DataNode接收流的代码关键在DataXceiver.run()中:

// DataXceiver.java
public void run() {
    Op op = Op.readOp(in);
    switch (op) {
        case WRITE_BLOCK:
            writeBlock(pipelineInfo, header, ...);
            break;
    }
}
private void writeBlock(...) throws IOException {
    // 1. 创建块接收流
    BlockReceiver blockReceiver = new BlockReceiver(block, storageType, in,
        sock.getRemoteSocketAddress(), datanode, diskBalancerEnabled, ...);
    // 2. 接收数据,同时写入本地和转发给下一个datanode
    try {
        blockReceiver.receiveBlock(mirrorOut, mirrorAddr, dataOut, ...);
    } finally {
        blockReceiver.close();
    }
}

4.5 集群配置示例(hdfs-site.xml)

<?xml version="1.0"?>
<configuration>
    <property>
        <name>dfs.nameservices</name>
        <value>mycluster</value>
    </property>
    <property>
        <name>dfs.ha.namenodes.mycluster</name>
        <value>nn1,nn2</value>
    </property>
    <property>
        <name>dfs.namenode.rpc-address.mycluster.nn1</name>
        <value>host1:8020</value>
    </property>
    <property>
        <name>dfs.namenode.rpc-address.mycluster.nn2</name>
        <value>host2:8020</value>
    </property>
    <property>
        <name>dfs.namenode.shared.edits.dir</name>
        <value>qjournal://host3:8485;host4:8485;host5:8485/mycluster</value>
    </property>
    <property>
        <name>dfs.client.failover.proxy.provider.mycluster</name>
        <value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value>
    </property>
</configuration>

五、效果数据:三种架构压测对比

测试环境:

  • 硬件:50台虚拟机,每台8 vCPU / 32GB RAM / 1TB HDD,万兆网卡
  • Hadoop版本:Apache Hadoop 3.2.1
  • 客户端:10台同配置机器并发写入
  • 数据:100个100MB文件(共10GB),块大小128MB,3副本
  • NameNode配置:堆内存16GB,GC算法G1,RPC handler数64

压测工具为dfsio,测试均为3次取平均。

架构 写入吞吐 (MB/s) 写入延迟 P99 (ms) NameNode CPU利用率 NameNode GC停顿(P99)
标准单NameNode 820 240 78% 1.2s
HA (Active/Standby) 795 260 82% 1.5s
Federation (2个NameNode) 1510 130 平均45% 0.4s

结论:HA因同步EditLog增加10~20ms延迟,吞吐略降;Federation吞吐近线性扩展,但文件分布在两个namespace,客户端需感知路由。

六、避坑指南:实际踩过的坑

在多次集群故障与优化中,我总结了以下HDFS相关坑,每一个都曾导致线上问题。

坑1:NameNode堆内存配置过大导致Full GC频繁

我曾将NN堆内存设为48GB(G1),结果每次Full GC超过5秒。原因是G1对超大堆的region数处理不当,且元数据对象引用链复杂。推荐:NN堆内存不超过32GB,若元数据量巨大(超过10亿对象),改用Federation分散。GC算法用CMS+G1均可,但务必开启ParallelRefProcEnabled,减少引用处理停顿。

坑2:FsImage合并(Checkpoint)时OOM

Standby NameNode或Secondary NameNode在做Checkpoint时,需要加载当前FsImage并与EditLog合并,如果EditLog积累过多(如100万+条),合并过程会创建大量临时对象导致OOM。解决:缩短EditLog滚动时间(dfs.namenode.checkpoint.period设为60秒),并调大dfs.namenode.checkpoint.txns为50万条。同时给Checkpoint的NN分配足够堆内存。

坑3:DataNode磁盘不均衡导致写入倾斜

新加入磁盘的DataNode会被HDFS自动感知,但旧DataNode磁盘使用率可能高达95%以上。HDFS默认均衡带宽只有20MB/s,对于多TB场景太慢。我踩坑:某次扩容后一个月仍有30%磁盘不均衡。必须手动调整dfs.datanode.balance.bandwidthPerSec到200MB/s,并定期执行hdfs balancer -threshold 5

坑4:小文件问题:元数据爆炸与RPC压力

一个Hive表有500万个文件(每个几KB),NameNode内存占用近800MB,且RPC耗时增加10倍。解决:使用HBase/Parquet + 适当合并小文件,或开启HDFS的dfs.namenode.file-close-num-committed-allowed缓存。切忌直接存储海量小文件。

七、总结(只一句话)

理解HDFS源码的核心在于NameNode的元数据管理和DataNode的流式写入,选择HA还是Federation取决于你的元数据规模与可用性要求,而避坑的关键是谨慎配置内存与GC。

<<>>