Hadoop 3.x 高可用集群 — 知识点详解
一、HDFS 高可用集群
1.1 HDFS HA 架构概述
核心知识点:
在 Hadoop 1.x 中,NameNode 存在单点故障(SPOF)。HDFS HA 通过配置 Active/Standby 两个 NameNode 来解决此问题。
关键组件:
组件作用
| Active NameNode |
处理所有客户端请求,维护文件系统元数据 |
| Standby NameNode |
作为热备,同步 Active 的编辑日志,随时接管 |
| JournalNode (JN) |
共享存储系统,存储 EditLog,保证两个 NN 之间数据同步 |
| ZKFailoverController (ZKFC) |
监控 NameNode 健康状态,通过 ZooKeeper 实现自动故障转移 |
| ZooKeeper |
提供分布式协调服务,维护 Active/Standby 的锁(ephemeral node) |
| DataNode |
向两个 NameNode 同时发送 Block 报告和心跳 |
架构图文字描述:
┌──────────────┐
│ ZooKeeper │
│ Cluster │
└──┬───────┬───┘
│ │
ZKFC │ │ ZKFC
┌──────┴──┐ ┌──┴──────┐
│ Active │ │ Standby │
│ NameNode│ │ NameNode│
└────┬────┘ └────┬────┘
│ │
┌────┴───────────┴────┐
│ JournalNode 集群 │
│ (至少3个,奇数个) │
└─────────────────────┘
│ │
┌──────────┴───────────┴──────────┐
│ DataNode 1, 2, 3 … N │
└─────────────────────────────────┘
1.2 JournalNode 工作机制
知识点:
- JournalNode 是一个轻量级的守护进程,通常部署奇数个(至少3个)
- Active NameNode 将 EditLog 写入 JournalNode 集群
- Standby NameNode 从 JournalNode 集群读取 EditLog 并应用到内存
- JournalNode 使用 Paxos 协议 保证数据一致性,需要超过半数(N/2+1)节点写入成功
配置 JournalNode 的 hdfs-site.xml 核心参数:
<property>
<name>dfs.namenode.shared.edits.dir</name>
<value>qjournal://node1:8485;node2:8485;node3:8485/mycluster</value>
</property>
<property>
<name>dfs.journalnode.edits.dir</name>
<value>/opt/module/hadoop-3.1.3/data/journalnode</value>
</property>
<property>
<name>dfs.journalnode.rpc-address</name>
<value>0.0.0.0:8485</value>
</property>
<property>
<name>dfs.journalnode.http-address</name>
<value>0.0.0.0:8480</value>
</property>
1.3 HDFS Federation(联邦)与 HA 的区别
知识点:
对比项HDFS FederationHDFS HA
| 目的 |
解决单个 NameNode 内存瓶颈 |
解决 NameNode 单点故障 |
| NameNode 数量 |
多个 NN 管理不同的命名空间 |
两个 NN(Active + Standby)管理同一个命名空间 |
| 元数据隔离 |
不同 NN 管理不同的 BlockPool |
共享同一份元数据 |
| 是否互补 |
可以与 HA 结合使用 |
可以与 Federation 结合使用 |
Federation 配置示例:
<property>
<name>dfs.nameservices</name>
<value>ns1,ns2</value>
</property>
<property>
<name>dfs.namenode.rpc-address.ns1</name>
<value>node1:8020</value>
</property>
<property>
<name>dfs.namenode.rpc-address.ns2</name>
<value>node2:8020</value>
</property>
1.4 HDFS HA 完整配置详解
1.4.1 core-site.xml 配置
<property>
<name>fs.defaultFS</name>
<value>hdfs://mycluster</value>
</property>
<property>
<name>ha.zookeeper.quorum</name>
<value>node1:2181,node2:2181,node3:2181</value>
</property>
<property>
<name>hadoop.tmp.dir</name>
<value>/opt/module/hadoop-3.1.3/data</value>
</property>
1.4.2 hdfs-site.xml 完整 HA 配置
<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>node1:8020</value>
</property>
<property>
<name>dfs.namenode.rpc-address.mycluster.nn2</name>
<value>node2:8020</value>
</property>
<property>
<name>dfs.namenode.http-address.mycluster.nn1</name>
<value>node1:9870</value>
</property>
<property>
<name>dfs.namenode.http-address.mycluster.nn2</name>
<value>node2:9870</value>
</property>
<property>
<name>dfs.namenode.shared.edits.dir</name>
<value>qjournal://node1:8485;node2:8485;node3:8485/mycluster</value>
</property>
<property>
<name>dfs.journalnode.edits.dir</name>
<value>/opt/module/hadoop-3.1.3/data/journalnode</value>
</property>
<property>
<name>dfs.client.failover.proxy.provider.mycluster</name>
<value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value>
</property>
<property>
<name>dfs.ha.fencing.methods</name>
<value>
sshfence
shell(/bin/true)
</value>
</property>
<property>
<name>dfs.ha.fencing.ssh.private-key-files</name>
<value>/home/hadoop/.ssh/id_rsa</value>
</property>
<property>
<name>dfs.ha.fencing.ssh.connect-timeout</name>
<value>30000</value>
</property>
<property>
<name>dfs.ha.automatic-failover.enabled</name>
<value>true</value>
</property>
<property>
<name>ha.zookeeper.session-timeout.ms</name>
<value>10000</value>
</property>
1.5 自动故障转移流程
知识点:
正常状态:
nn1 (Active) 持有 ZooKeeper 的 ephemeral node(临时节点)锁
nn2 (Standby) 尝试获取锁但失败,保持 Standby
故障发生:
1. nn1 进程崩溃或网络断开
2. nn1 的 ZKFC 与 ZooKeeper 的 session 超时
3. ZooKeeper 删除 nn1 的 ephemeral node
4. nn2 的 ZKFC 发现锁被释放,立即创建自己的 ephemeral node
5. nn2 的 ZKFC 获取到锁,调用 nn2 变为 Active
6. nn2 从 JournalNode 集群读取所有 EditLog 并应用(元数据与 nn1 同步)
7. nn2 开始对外提供服务
注意:旧 nn1 恢复后自动变为 Standby
手动故障转移命令:
hdfs haadmin -getServiceState nn1
hdfs haadmin -getServiceState nn2
hdfs haadmin -transitionToActive nn1
hdfs haadmin -transitionToStandby nn2
hdfs haadmin -failover nn1 nn2
hdfs haadmin -failover nn1 nn2 –forcefence
hdfs haadmin -failover nn1 nn2 –forceactive
1.6 HDFS HA 启动顺序(手动方式)
zkServer.sh start
zkServer.sh status
hdfs –daemon start journalnode
hdfs namenode -format
hdfs –daemon start namenode
hdfs namenode -bootstrapStandby
hdfs –daemon start namenode
hdfs zkfc -formatZK
hdfs –daemon start datanode
hdfs –daemon start zkfc
hdfs haadmin -getServiceState nn1
hdfs haadmin -getServiceState nn2
一键启动脚本(基于 start-dfs.sh):
#!/bin/bash
echo "============ 1. 启动 ZooKeeper 集群 ============"
for host in node1 node2 node3; do
echo "在 $host 上启动 ZooKeeper…"
ssh $host "/opt/module/zookeeper-3.5.7/bin/zkServer.sh start"
done
sleep 5
echo "============ 2. 启动 JournalNode ============"
for host in node1 node2 node3; do
echo "在 $host 上启动 JournalNode…"
ssh $host "/opt/module/hadoop-3.1.3/bin/hdfs –daemon start journalnode"
done
sleep 3
echo "============ 3. 启动 HDFS(含 NameNode、DataNode、ZKFC) ============"
/opt/module/hadoop-3.1.3/sbin/start-dfs.sh
echo "============ 4. 验证集群状态 ============"
hdfs dfsadmin -report
echo "nn1 状态: $(hdfs haadmin -getServiceState nn1)"
echo "nn2 状态: $(hdfs haadmin -getServiceState nn2)"
echo "HDFS HA 集群启动完成!"
1.7 HDFS HA Java 客户端代码
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.*;
import org.apache.hadoop.io.IOUtils;
import java.io.*;
import java.net.URI;
public class HdfsHAClient {
private FileSystem fileSystem;
public void init() throws Exception {
Configuration conf = new Configuration();
conf.set("fs.defaultFS", "hdfs://mycluster");
conf.set("ha.zookeeper.quorum", "node1:2181,node2:2181,node3:2181");
fileSystem = FileSystem.get(new URI("hdfs://mycluster"), conf, "hadoop");
}
public void uploadFile(String localPath, String hdfsPath) throws Exception {
Path srcPath = new Path(localPath);
Path dstPath = new Path(hdfsPath);
fileSystem.copyFromLocalFile(false, true, srcPath, dstPath);
System.out.println("文件上传成功: " + localPath + " -> " + hdfsPath);
}
public void downloadFile(String hdfsPath, String localPath) throws Exception {
Path srcPath = new Path(hdfsPath);
Path dstPath = new Path(localPath);
fileSystem.copyToLocalFile(false, srcPath, dstPath);
System.out.println("文件下载成功: " + hdfsPath + " -> " + localPath);
}
public void writeFileByStream(String hdfsPath, String content) throws Exception {
Path path = new Path(hdfsPath);
FSDataOutputStream outputStream = fileSystem.create(path);
outputStream.writeBytes(content);
outputStream.hflush();
outputStream.close();
System.out.println("流式写入成功: " + hdfsPath);
}
public String readFileByStream(String hdfsPath) throws Exception {
Path path = new Path(hdfsPath);
if (!fileSystem.exists(path)) {
System.out.println("文件不存在: " + hdfsPath);
return null;
}
FSDataInputStream inputStream = fileSystem.open(path);
StringBuilder content = new StringBuilder();
String line;
while ((line = inputStream.readLine()) != null) {
content.append(line).append("\\n");
}
inputStream.close();
return content.toString();
}
public void listDirectory(String dirPath) throws Exception {
Path path = new Path(dirPath);
if (!fileSystem.exists(path)) {
System.out.println("目录不存在: " + dirPath);
return;
}
FileStatus[] fileStatuses = fileSystem.listStatus(path);
System.out.println("======== 目录内容: " + dirPath + " ========");
for (FileStatus status : fileStatuses) {
String type = status.isDirectory() ? "[目录]" : "[文件]";
String name = status.getPath().getName();
long size = status.getLen();
String permission = status.getPermission().toString();
System.out.printf("%-8s %-30s 大小: %-10d 权限: %s%n",
type, name, size, permission);
}
}
public void listFilesRecursive(String dirPath) throws Exception {
Path path = new Path(dirPath);
RemoteIterator<LocatedFileStatus> fileIterator = fileSystem.listFiles(path, true);
System.out.println("======== 递归文件列表 ========");
while (fileIterator.hasNext()) {
LocatedFileStatus fileStatus = fileIterator.next();
String filePath = fileStatus.getPath().toString();
long size = fileStatus.getLen();
BlockLocation[] blockLocations = fileStatus.getBlockLocations();
System.out.printf("文件: %-60s 大小: %-10d 块数: %d%n",
filePath, size, blockLocations.length);
}
}
public void deleteFile(String path, boolean recursive) throws Exception {
Path filePath = new Path(path);
boolean result = fileSystem.delete(filePath, recursive);
if (result) {
System.out.println("删除成功: " + path);
} else {
System.out.println("删除失败: " + path);
}
}
public void mkdir(String dirPath) throws Exception {
Path path = new Path(dirPath);
boolean result = fileSystem.mkdirs(path);
if (result) {
System.out.println("目录创建成功: " + dirPath);
} else {
System.out.println("目录创建失败: " + dirPath);
}
}
public void close() throws Exception {
if (fileSystem != null) {
fileSystem.close();
System.out.println("HDFS 连接已关闭");
}
}
public static void main(String[] args) {
HdfsHAClient client = new HdfsHAClient();
try {
client.init();
client.mkdir("/ha-test");
client.writeFileByStream("/ha-test/hello.txt",
"Hello HDFS HA Cluster!\\nThis is a test file.\\n");
client.uploadFile("/tmp/local-file.txt", "/ha-test/uploaded.txt");
client.listDirectory("/ha-test");
String content = client.readFileByStream("/ha-test/hello.txt");
System.out.println("======== 文件内容 ========");
System.out.println(content);
client.listFilesRecursive("/");
client.downloadFile("/ha-test/hello.txt", "/tmp/downloaded.txt");
client.deleteFile("/ha-test/hello.txt", false);
client.deleteFile("/ha-test", true);
} catch (Exception e) {
e.printStackTrace();
} finally {
try {
client.close();
} catch (Exception e) {
e.printStackTrace();
}
}
}
}
1.8 HDFS HA 状态检查与故障排查脚本
#!/bin/bash
echo "============ HDFS HA 集群健康检查 ============"
echo "检查时间: $(date '+%Y-%m-%d %H:%M:%S')"
echo ""
echo "【1】ZooKeeper 集群状态:"
for host in node1 node2 node3; do
status=$(ssh $host "/opt/module/zookeeper-3.5.7/bin/zkServer.sh status 2>&1")
echo " $host: $status"
done
echo ""
echo "【2】JournalNode 进程状态:"
for host in node1 node2 node3; do
jn_pid=$(ssh $host "jps -l | grep JournalNode")
if [ -n "$jn_pid" ]; then
echo " $host: 运行中 – $jn_pid"
else
echo " $host: 未运行!"
fi
done
echo ""
echo "【3】NameNode HA 状态:"
nn1_state=$(hdfs haadmin -getServiceState nn1 2>/dev/null)
nn2_state=$(hdfs haadmin -getServiceState nn2 2>/dev/null)
echo " nn1 (node1): ${nn1_state:-'无法连接'}"
echo " nn2 (node2): ${nn2_state:-'无法连接'}"
if [ "$nn1_state" = "active" ] && [ "$nn2_state" = "standby" ]; then
echo " HA 状态: 正常 (nn1=Active, nn2=Standby)"
elif [ "$nn1_state" = "standby" ] && [ "$nn2_state" = "active" ]; then
echo " HA 状态: 正常 (nn1=Standby, nn2=Active)"
else
echo " HA 状态: 异常!请检查 NameNode"
fi
echo ""
echo "【4】DataNode 状态:"
live_nodes=$(hdfs dfsadmin -report 2>/dev/null | grep "Live datanodes" | grep -oP '\\d+')
dead_nodes=$(hdfs dfsadmin -report 2>/dev/null | grep "Dead datanodes" | grep -oP '\\d+')
echo " 活跃 DataNode: ${live_nodes:-0}"
echo " 死亡 DataNode: ${dead_nodes:-0}"
echo ""
echo "【5】HDFS 读写测试:"
test_dir="/ha-health-check-$(date +%s)"
if hdfs dfs -mkdir -p $test_dir 2>/dev/null; then
echo " 写入测试: 成功"
if hdfs dfs -ls $test_dir 2>/dev/null; then
echo " 读取测试: 成功"
else
echo " 读取测试: 失败"
fi
hdfs dfs -rm -r $test_dir 2>/dev/null
else
echo " 写入测试: 失败!HDFS 可能不可用"
fi
echo ""
echo "============ 检查完成 ============"
二、YARN 高可用集群
2.1 YARN HA 架构概述
核心知识点:
YARN HA 与 HDFS HA 类似,通过配置 Active/Standby 两个 ResourceManager 来解决单点故障问题。
关键组件:
组件作用
| Active ResourceManager |
处理客户端请求,管理资源分配,调度 Application |
| Standby ResourceManager |
热备,同步 Active RM 的状态,随时接管 |
| ZKFailoverController (RM ZKFC) |
内嵌在 RM 中,负责与 ZooKeeper 交互实现自动故障转移 |
| ZooKeeper |
存储 RM 的选举信息,维护 Active/Standby 的锁 |
| ResourceManagerStateStore |
存储 RM 的应用状态(内存或 ZooKeeper) |
| NodeManager |
向所有 RM 注册,但只接收 Active RM 的指令 |
YARN HA 与 HDFS HA 的对比:
对比项HDFS HAYARN HA
| 共享存储 |
JournalNode |
ZooKeeper(RMStateStore) |
| 状态同步方式 |
EditLog 实时同步 |
应用状态存储到 ZK |
| ZKFC |
独立进程 |
内嵌在 ResourceManager 中 |
| 故障转移影响 |
客户端透明切换 |
正在运行的 Application 需要重新提交 |
2.2 YARN HA 完整配置
2.2.1 yarn-site.xml 配置
<property>
<name>yarn.resourcemanager.ha.enabled</name>
<value>true</value>
</property>
<property>
<name>yarn.resourcemanager.cluster-id</name>
<value>yarn-cluster</value>
</property>
<property>
<name>yarn.resourcemanager.ha.rm-ids</name>
<value>rm1,rm2</value>
</property>
<property>
<name>yarn.resourcemanager.hostname.rm1</name>
<value>node1</value>
</property>
<property>
<name>yarn.resourcemanager.hostname.rm2</name>
<value>node2</value>
</property>
<property>
<name>yarn.resourcemanager.webapp.address.rm1</name>
<value>node1:8088</value>
</property>
<property>
<name>yarn.resourcemanager.webapp.address.rm2</name>
<value>node2:8088</value>
</property>
<property>
<name>yarn.resourcemanager.address.rm1</name>
<value>node1:8032</value>
</property>
<property>
<name>yarn.resourcemanager.address.rm2</name>
<value>node2:8032</value>
</property>
<property>
<name>yarn.resourcemanager.scheduler.address.rm1</name>
<value>node1:8030</value>
</property>
<property>
<name>yarn.resourcemanager.scheduler.address.rm2</name>
<value>node2:8030</value>
</property>
<property>
<name>yarn.resourcemanager.resource-tracker.address.rm1</name>
<value>node1:8031</value>
</property>
<property>
<name>yarn.resourcemanager.resource-tracker.address.rm2</name>
<value>node2:8031</value>
</property>
<property>
<name>yarn.resourcemanager.admin.address.rm1</name>
<value>node1:8033</value>
</property>
<property>
<name>yarn.resourcemanager.admin.address.rm2</name>
<value>node2:8033</value>
</property>
<property>
<name>ha.zookeeper.quorum</name>
<value>node1:2181,node2:2181,node3:2181</value>
</property>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.enabled</name>
<value>true</value>
</property>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.embedded</name>
<value>true</value>
</property>
可选值:
– org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore
将状态存储在 ZooKeeper 中(推荐,适合 HA)
– org.apache.hadoop.yarn.server.resourcemanager.recovery.FileSystemRMStateStore
将状态存储在 HDFS 中
– org.apache.hadoop.yarn.server.resourcemanager.recovery.LeveldbRMStateStore
将状态存储在 LevelDB 中
–>
<property>
<name>yarn.resourcemanager.store.class</name>
<value>org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore</value>
</property>
<property>
<name>yarn.nodemanager.aux-services</name>
<value>mapreduce_shuffle</value>
</property>
<property>
<name>yarn.nodemanager.resource.memory-mb</name>
<value>4096</value>
</property>
<property>
<name>yarn.nodemanager.resource.cpu-vcores</name>
<value>4</value>
</property>
<property>
<name>yarn.resourcemanager.scheduler.class</name>
<value>org.apache.hadoop.yarn.server.resourcemanager.scheduler.capacity.CapacityScheduler</value>
</property>
<property>
<name>yarn.log-aggregation-enable</name>
<value>true</value>
</property>
<property>
<name>yarn.log-aggregation.retain-seconds</name>
<value>604800</value>
</property>
2.3 YARN HA 故障转移流程
正常状态:
rm1 (Active) 在 ZooKeeper 上创建 ephemeral node
rm2 (Standby) 监听 ZK 上的节点变化
所有 NodeManager 向两个 RM 都注册,但只接受 Active RM 的指令
故障发生:
1. rm1 进程崩溃
2. rm1 在 ZooKeeper 上的 ephemeral node 自动删除
3. rm2 的 StandbyElector 检测到节点删除事件
4. rm2 创建自己的 ephemeral node,成为新的 Active
5. rm2 从 ZooKeeper 的 ZKRMStateStore 中恢复应用状态
6. rm2 向所有 NodeManager 重新注册
7. 正在运行的 Application 会经历短暂中断后恢复
8. 已提交但未运行的 Application 由新 Active RM 重新调度
注意:
– Container 级别的任务如果正在运行,通常可以继续执行
– ApplicationMaster 需要重新与新 Active RM 建立连接
– 配置了 retry 的客户端会自动重试连接到新 RM
2.4 YARN HA 启动与管理命令
/opt/module/hadoop-3.1.3/sbin/start-yarn.sh
yarn –daemon start resourcemanager
yarn –daemon start resourcemanager
yarn –daemon start nodemanager
yarn rmadmin -getServiceState rm1
yarn rmadmin -getServiceState rm2
yarn rmadmin -failover rm1 rm2
yarn rmadmin -transitionToActive rm1
yarn rmadmin -transitionToStandby rm2
yarn rmadmin -getAllServiceState
hadoop jar /opt/module/hadoop-3.1.3/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.1.3.jar \\
wordcount \\
/input \\
/output
yarn application -list
yarn application -status application_1234567890123_0001
yarn application -kill application_1234567890123_0001
2.5 YARN HA Java 客户端代码
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.yarn.api.records.*;
import org.apache.hadoop.yarn.client.api.YarnClient;
import org.apache.hadoop.yarn.exceptions.YarnException;
import java.io.IOException;
import java.util.List;
import java.util.EnumSet;
public class YarnHAClient {
private YarnClient yarnClient;
public void init() {
Configuration conf = new Configuration();
conf.setBoolean("yarn.resourcemanager.ha.enabled", true);
conf.set("yarn.resourcemanager.cluster-id", "yarn-cluster");
conf.set("yarn.resourcemanager.ha.rm-ids", "rm1,rm2");
conf.set("yarn.resourcemanager.hostname.rm1", "node1");
conf.set("yarn.resourcemanager.hostname.rm2", "node2");
conf.set("ha.zookeeper.quorum", "node1:2181,node2:2181,node3:2181");
yarnClient = YarnClient.createYarnClient();
yarnClient.init(conf);
yarnClient.start();
}
public void getClusterInfo() throws IOException, YarnException {
YarnClusterMetrics clusterMetrics = yarnClient.getYarnClusterMetrics();
System.out.println("============ YARN HA 集群信息 ============");
System.out.println("总节点数: " + clusterMetrics.getNumNodeManagers());
System.out.println("活跃节点数: " + clusterMetrics.getNumActiveNodeManagers());
System.out.println("不健康节点数: " + clusterMetrics.getUnhealthyNodeManagers());
System.out.println("已停用节点数: " + clusterMetrics.getNumDecommissionedNodeManagers());
System.out.println("丢失节点数: " + clusterMetrics.getNumLostNodeManagers());
}
public void listNodeManagers() throws IOException, YarnException {
List<NodeReport> nodeReports = yarnClient.getNodeReports(
EnumSet.of(NodeState.RUNNING));
System.out.println("============ NodeManager 详细信息 ============");
for (NodeReport node : nodeReports) {
String nodeId = node.getNodeId().toString();
String httpAddress = node.getHttpAddress();
long memoryTotal = node.getCapability().getMemorySize();
long memoryUsed = node.getUsed().getMemorySize();
int vcoresTotal = node.getCapability().getVirtualCores();
int vcoresUsed = node.getUsed().getVirtualCores();
int numContainers = node.getNumContainers();
String healthReport = node.getHealthReport();
System.out.printf("节点: %-20s HTTP: %-25s%n", nodeId, httpAddress);
System.out.printf(" 内存: %dMB / %dMB (已用/总量)%n", memoryUsed, memoryTotal);
System.out.printf(" CPU: %d / %d 核 (已用/总量)%n", vcoresUsed, vcoresTotal);
System.out.printf(" 容器数: %d%n", numContainers);
System.out.printf(" 健康报告: %s%n", healthReport.isEmpty() ? "正常" : healthReport);
System.out.println();
}
}
public void listApplications() throws IOException, YarnException {
List<ApplicationReport> apps = yarnClient.getApplications();
System.out.println("============ YARN 应用列表 ============");
System.out.printf("%-40s %-15s %-12s %-20s%n",
"Application ID", "Application Name", "State", "Tracking URL");
for (ApplicationReport app : apps) {
ApplicationId appId = app.getApplicationId();
String appName = app.getName();
YarnApplicationState state = app.getYarnApplicationState();
String trackingUrl = app.getTrackingUrl();
System.out.printf("%-40s %-15s %-12s %-20s%n",
appId.toString(), appName, state, trackingUrl);
}
}
public void getApplicationDetail(String applicationId)
throws IOException, YarnException {
ApplicationId appId = ApplicationId.fromString(applicationId);
ApplicationReport report = yarnClient.getApplicationReport(appId);
System.out.println("============ 应用详细信息 ============");
System.out.println("应用 ID: " + report.getApplicationId());
System.out.println("应用名称: " + report.getName());
System.out.println("应用类型: " + report.getApplicationType());
System.out.println("应用状态: " + report.getYarnApplicationState());
System.out.println("最终状态: " + report.getFinalApplicationStatus());
System.out.println("提交用户: " + report.getUser());
System.out.println("所在队列: " + report.getQueue());
System.out.println("提交时间: " + new java.util.Date(report.getSubmitTime()));
System.out.println("启动时间: " + new java.util.Date(report.getStartTime()));
long finishTime = report.getFinishTime();
if (finishTime > 0) {
System.out.println("结束时间: " + new java.util.Date(finishTime));
}
System.out.println("跟踪 URL: " + report.getTrackingUrl());
System.out.println("诊断信息: " + report.getDiagnostics());
}
public void listQueues() throws IOException, YarnException {
QueueInfo rootQueue = yarnClient.getQueueInfo("root");
System.out.println("============ YARN 队列信息 ============");
printQueueInfo(rootQueue, 0);
}
private void printQueueInfo(QueueInfo queue, int indent) {
String indentStr = " ".repeat(indent);
System.out.printf("%s队列名: %s%n", indentStr, queue.getQueueName());
System.out.printf("%s 状态: %s%n", indentStr, queue.getQueueState());
System.out.printf("%s 容量: %.1f%%%n", indentStr, queue.getCapacity() * 100);
System.out.printf("%s 最大容量: %.1f%%%n", indentStr, queue.getMaximumCapacity() * 100);
System.out.printf("%s 当前使用: %.1f%%%n", indentStr, queue.getCurrentCapacity() * 100);
System.out.printf("%s 应用数: %d%n", indentStr, queue.getApplications().size());
System.out.println();
List<QueueInfo> childQueues = queue.getChildQueues();
if (childQueues != null) {
for (QueueInfo child : childQueues) {
printQueueInfo(child, indent + 1);
}
}
}
public void close() {
if (yarnClient != null) {
yarnClient.stop();
System.out.println("YARN 客户端已关闭");
}
}
public static void main(String[] args) {
YarnHAClient client = new YarnHAClient();
try {
client.init();
client.getClusterInfo();
client.listNodeManagers();
client.listApplications();
client.listQueues();
} catch (Exception e) {
e.printStackTrace();
} finally {
client.close();
}
}
}
三、部署 Hadoop 高可用集群
3.1 环境规划
集群节点规划:
主机名IP 地址角色
| node1 |
192.168.10.101 |
NameNode(active), DataNode, ResourceManager(active), NodeManager, JournalNode, ZooKeeper, ZKFC |
| node2 |
192.168.10.102 |
NameNode(standby), DataNode, ResourceManager(standby), NodeManager, JournalNode, ZooKeeper, ZKFC |
| node3 |
192.168.10.103 |
DataNode, NodeManager, JournalNode, ZooKeeper |
软件版本规划:
软件版本
| JDK |
1.8 (jdk-8u212) |
| Hadoop |
3.1.3 |
| ZooKeeper |
3.5.7 |
端口规划:
服务端口说明
| NameNode RPC |
8020 |
客户端与 NN 的 RPC 通信 |
| NameNode HTTP |
9870 |
NN Web UI |
| DataNode |
9866 |
DN 数据传输端口 |
| DataNode HTTP |
9864 |
DN Web UI |
| ResourceManager HTTP |
8088 |
RM Web UI |
| ResourceManager RPC |
8032 |
客户端与 RM 的 RPC 通信 |
| NodeManager HTTP |
8042 |
NM Web UI |
| JournalNode RPC |
8485 |
JN RPC 通信端口 |
| JournalNode HTTP |
8480 |
JN Web UI |
| ZooKeeper |
2181 |
ZK 客户端连接端口 |
| ZooKeeper Leader Election |
2888, 3888 |
ZK 集群内部通信 |
3.2 基础环境准备
#!/bin/bash
sudo systemctl stop firewalld
sudo systemctl disable firewalld
echo "防火墙已关闭"
sudo hostnamectl set-hostname node1
echo "主机名已设置为 node1"
cat >> /etc/hosts << 'EOF'
192.168.10.101 node1
192.168.10.102 node2
192.168.10.103 node3
EOF
echo "hosts 文件已配置"
ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa
for host in node1 node2 node3; do
ssh-copy-id $host
echo "已配置到 $host 的免密登录"
done
tar -zxvf /opt/software/jdk-8u212-linux-x64.tar.gz -C /opt/module/
cat >> ~/.bash_profile << 'EOF'
# Java Environment
export JAVA_HOME=/opt/module/jdk1.8.0_212
export PATH=$PATH:$JAVA_HOME/bin
export CLASSPATH=.:$JAVA_HOME/lib/dt.jar:$JAVA_HOME/lib/tools.jar
EOF
source ~/.bash_profile
java -version
echo "JDK 安装完成"
sudo yum install -y chrony
sudo systemctl start chronyd
sudo systemctl enable chronyd
echo "时间同步服务已配置"
3.3 ZooKeeper 集群部署
tar -zxvf /opt/software/apache-zookeeper-3.5.7-bin.tar.gz -C /opt/module/
mv /opt/module/apache-zookeeper-3.5.7-bin /opt/module/zookeeper-3.5.7
cat >> ~/.bash_profile << 'EOF'
# ZooKeeper Environment
export ZOOKEEPER_HOME=/opt/module/zookeeper-3.5.7
export PATH=$PATH:$ZOOKEEPER_HOME/bin
EOF
source ~/.bash_profile
cp $ZOOKEEPER_HOME/conf/zoo_sample.cfg $ZOOKEEPER_HOME/conf/zoo.cfg
cat > $ZOOKEEPER_HOME/conf/zoo.cfg << 'EOF'
# ZooKeeper 服务器之间或客户端与服务器之间的心跳间隔(毫秒)
# 每隔 2000ms 发送一次心跳
tickTime=2000
# Follower 服务器初始连接到 Leader 时的最大心跳数
# 即初始化连接时最长能忍受 tickTime * initLimit = 2000 * 10 = 20000ms = 20秒
initLimit=10
# Follower 服务器与 Leader 服务器之间请求和应答的最大心跳数
# 即通信超时时长为 tickTime * syncLimit = 2000 * 5 = 10000ms = 10秒
syncLimit=5
# ZooKeeper 数据存储目录(存放内存数据快照和事务日志)
dataDir=/opt/module/zookeeper-3.5.7/zkData
# ZooKeeper 客户端连接端口
clientPort=2181
# 集群服务器配置
# server.A=B:C:D
# A: 服务器编号(对应 myid 文件中的数字)
# B: 服务器的 IP 地址或主机名
# C: Follower 与 Leader 交换信息的端口(数据同步端口)
# D: 选举端口(Leader 选举时使用的端口)
server.1=node1:2888:3888
server.2=node2:2888:3888
server.3=node3:2888:3888
# 配置 ZooKeeper 的 4 字命令(白名单)
# 包括 stat, ruok, conf, isro 等管理命令
# "四字命令"是指通过 telnet 或 nc 发送的 4 个字符命令
4lw.commands.whitelist=*
EOF
mkdir -p $ZOOKEEPER_HOME/zkData
echo "1" > $ZOOKEEPER_HOME/zkData/myid
echo "ZooKeeper 配置完成"
for host in node2 node3; do
rsync -az /opt/module/zookeeper-3.5.7 $host:/opt/module/
echo "已同步到 $host"
done
for host in node1 node2 node3; do
echo "在 $host 上启动 ZooKeeper…"
ssh $host "/opt/module/zookeeper-3.5.7/bin/zkServer.sh start"
done
sleep 5
for host in node1 node2 node3; do
echo "==== $host ===="
ssh $host "/opt/module/zookeeper-3.5.7/bin/zkServer.sh status"
done
3.4 Hadoop 安装与配置
tar -zxvf /opt/software/hadoop-3.1.3.tar.gz -C /opt/module/
cat >> ~/.bash_profile << 'EOF'
# Hadoop Environment
export HADOOP_HOME=/opt/module/hadoop-3.1.3
export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin
EOF
source ~/.bash_profile
cat >> $HADOOP_HOME/etc/hadoop/hadoop-env.sh << 'EOF'
# 指定 JDK 安装路径
export JAVA_HOME=/opt/module/jdk1.8.0_212
# 指定 HDFS 相关进程的运行用户
# 如果不配置,启动时可能会报 "the xxx is not allowed to run" 错误
export HDFS_NAMENODE_USER=hadoop
export HDFS_DATANODE_USER=hadoop
export HDFS_JOURNALNODE_USER=hadoop
export HDFS_ZKFC_USER=hadoop
export HDFS_SECONDARYNAMENODE_USER=hadoop
# 指定 YARN 相关进程的运行用户
export YARN_RESOURCEMANAGER_USER=hadoop
export YARN_NODEMANAGER_USER=hadoop
EOF
cat > $HADOOP_HOME/etc/hadoop/core-site.xml << 'EOF'
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!– 指定 HDFS 的默认文件系统名称为 nameservice 的逻辑名 –>
<property>
<name>fs.defaultFS</name>
<value>hdfs://mycluster</value>
</property>
<!– 指定 Hadoop 临时数据存储目录 –>
<property>
<name>hadoop.tmp.dir</name>
<value>/opt/module/hadoop-3.1.3/data</value>
</property>
<!– 指定 ZooKeeper 集群地址 –>
<property>
<name>ha.zookeeper.quorum</name>
<value>node1:2181,node2:2181,node3:2181</value>
</property>
</configuration>
EOF
cat > $HADOOP_HOME/etc/hadoop/hdfs-site.xml << 'EOF'
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!– NameService 逻辑名称 –>
<property>
<name>dfs.nameservices</name>
<value>mycluster</value>
</property>
<!– 两个 NameNode 的标识符 –>
<property>
<name>dfs.ha.namenodes.mycluster</name>
<value>nn1,nn2</value>
</property>
<!– nn1 RPC 地址 –>
<property>
<name>dfs.namenode.rpc-address.mycluster.nn1</name>
<value>node1:8020</value>
</property>
<!– nn2 RPC 地址 –>
<property>
<name>dfs.namenode.rpc-address.mycluster.nn2</name>
<value>node2:8020</value>
</property>
<!– nn1 HTTP 地址 –>
<property>
<name>dfs.namenode.http-address.mycluster.nn1</name>
<value>node1:9870</value>
</property>
<!– nn2 HTTP 地址 –>
<property>
<name>dfs.namenode.http-address.mycluster.nn2</name>
<value>node2:9870</value>
</property>
<!– JournalNode 集群地址 –>
<property>
<name>dfs.namenode.shared.edits.dir</name>
<value>qjournal://node1:8485;node2:8485;node3:8485/mycluster</value>
</property>
<!– JournalNode 数据存储目录 –>
<property>
<name>dfs.journalnode.edits.dir</name>
<value>/opt/module/hadoop-3.1.3/data/journalnode</value>
</property>
<!– 客户端故障转移代理类 –>
<property>
<name>dfs.client.failover.proxy.provider.mycluster</name>
<value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value>
</property>
<!– 隔离机制:防止脑裂 –>
<property>
<name>dfs.ha.fencing.methods</name>
<value>sshfence(shell(/bin/true))</value>
</property>
<!– SSH 私钥路径 –>
<property>
<name>dfs.ha.fencing.ssh.private-key-files</name>
<value>/home/hadoop/.ssh/id_rsa</value>
</property>
<!– 开启自动故障转移 –>
<property>
<name>dfs.ha.automatic-failover.enabled</name>
<value>true</value>
</property>
<!– 副本数 –>
<property>
<name>dfs.replication</name>
<value>3</value>
</property>
<!– 关闭权限检查(测试环境) –>
<property>
<name>dfs.permissions.enabled</name>
<value>false</value>
</property>
</configuration>
EOF
cat > $HADOOP_HOME/etc/hadoop/yarn-site.xml << 'EOF'
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!– 开启 ResourceManager HA –>
<property>
<name>yarn.resourcemanager.ha.enabled</name>
<value>true</value>
</property>
<!– 集群 ID –>
<property>
<name>yarn.resourcemanager.cluster-id</name>
<value>yarn-cluster</value>
</property>
<!– ResourceManager 标识符 –>
<property>
<name>yarn.resourcemanager.ha.rm-ids</name>
<value>rm1,rm2</value>
</property>
<!– rm1 主机名 –>
<property>
<name>yarn.resourcemanager.hostname.rm1</name>
<value>node1</value>
</property>
<!– rm2 主机名 –>
<property>
<name>yarn.resourcemanager.hostname.rm2</name>
<value>node2</value>
</property>
<!– rm1 Web UI –>
<property>
<name>yarn.resourcemanager.webapp.address.rm1</name>
<value>node1:8088</value>
</property>
<!– rm2 Web UI –>
<property>
<name>yarn.resourcemanager.webapp.address.rm2</name>
<value>node2:8088</value>
</property>
<!– ZooKeeper 地址 –>
<property>
<name>ha.zookeeper.quorum</name>
<value>node1:2181,node2:2181,node3:2181</value>
</property>
<!– 开启自动故障转移 –>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.enabled</name>
<value>true</value>
</property>
<!– 内嵌式选举 –>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.embedded</name>
<value>true</value>
</property>
<!– 状态存储在 ZooKeeper –>
<property>
<name>yarn.resourcemanager.store.class</name>
<value>org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore</value>
</property>
<!– NodeManager 辅助服务 –>
<property>
<name>yarn.nodemanager.aux-services</name>
<value>mapreduce_shuffle</value>
</property>
<!– 开启日志聚合 –>
<property>
<name>yarn.log-aggregation-enable</name>
<value>true</value>
</property>
<!– 日志保留 7 天 –>
<property>
<name>yarn.log-aggregation.retain-seconds</name>
<value>604800</value>
</property>
<!– NodeManager 可用内存 –>
<property>
<name>yarn.nodemanager.resource.memory-mb</name>
<value>4096</value>
</property>
<!– NodeManager 可用 CPU 核数 –>
<property>
<name>yarn.nodemanager.resource.cpu-vcores</name>
<value>4</value>
</property>
</configuration>
EOF
cat > $HADOOP_HOME/etc/hadoop/mapred-site.xml << 'EOF'
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!– 指定 MapReduce 运行在 YARN 上 –>
<!– 可选值: local(本地模式)、classic(经典模式)、yarn –>
<property>
<name>mapreduce.framework.name</name>
<value>yarn</value>
</property>
<!– 历史服务器地址(用于查看已完成的 MapReduce 作业历史) –>
<property>
<name>mapreduce.jobhistory.address</name>
<value>node1:10020</value>
</property>
<!– 历史服务器 Web UI 地址 –>
<property>
<name>mapreduce.jobhistory.webapp.address</name>
<value>node1:19888</value>
</property>
</configuration>
EOF
cat > $HADOOP_HOME/etc/hadoop/workers << 'EOF'
node1
node2
node3
EOF
for host in node2 node3; do
echo "正在同步 Hadoop 到 $host…"
rsync -az /opt/module/hadoop-3.1.3 $host:/opt/module/
rsync -az ~/.bash_profile $host:~/
echo "$host 同步完成"
done
echo "Hadoop 配置完成,已分发到所有节点"
3.5 集群初始化与启动
#!/bin/bash
echo "============ Hadoop HA 集群初始化与启动 ============"
echo "【步骤1】启动 ZooKeeper 集群…"
for host in node1 node2 node3; do
ssh $host "/opt/module/zookeeper-3.5.7/bin/zkServer.sh start"
done
sleep 5
for host in node1 node2 node3; do
echo "$host: $(ssh $host '/opt/module/zookeeper-3.5.7/bin/zkServer.sh status' | grep Mode)"
done
echo ""
echo "【步骤2】启动 JournalNode…"
for host in node1 node2 node3; do
ssh $host "/opt/module/hadoop-3.1.3/bin/hdfs –daemon start journalnode"
done
sleep 3
for host in node1 node2 node3; do
jn=$(ssh $host "jps | grep JournalNode")
echo "$host: $jn"
done
echo ""
echo "【步骤3】格式化 NameNode(仅在 node1 上执行)…"
/opt/module/hadoop-3.1.3/bin/hdfs namenode -format
echo ""
echo "【步骤4】启动 node1 的 NameNode…"
/opt/module/hadoop-3.1.3/bin/hdfs –daemon start namenode
sleep 3
echo "node1 NameNode 已启动"
echo ""
echo "【步骤5】在 node2 上同步 NameNode 元数据…"
ssh node2 "/opt/module/hadoop-3.1.3/bin/hdfs namenode -bootstrapStandby"
echo "启动 node2 的 NameNode…"
ssh node2 "/opt/module/hadoop-3.1.3/bin/hdfs –daemon start namenode"
sleep 3
echo "node2 NameNode 已启动"
echo ""
echo "【步骤6】格式化 ZKFC…"
/opt/module/hadoop-3.1.3/bin/hdfs zkfc -formatZK
echo ""
echo "【步骤7】启动 HDFS 集群…"
/opt/module/hadoop-3.1.3/sbin/start-dfs.sh
echo ""
echo "【步骤8】启动 YARN 集群…"
/opt/module/hadoop-3.1.3/sbin/start-yarn.sh
echo ""
echo "【步骤9】启动 MapReduce 历史服务器…"
ssh node1 "/opt/module/hadoop-3.1.3/bin/mapred –daemon start historyserver"
echo ""
echo "============ 集群状态验证 ============"
sleep 10
echo "HDFS HA 状态:"
echo " nn1: $(hdfs haadmin -getServiceState nn1)"
echo " nn2: $(hdfs haadmin -getServiceState nn2)"
echo "YARN HA 状态:"
echo " rm1: $(yarn rmadmin -getServiceState rm1)"
echo " rm2: $(yarn rmadmin -getServiceState rm2)"
echo ""
echo "HDFS 集群报告:"
hdfs dfsadmin -report | head -20
echo ""
echo "各节点 Java 进程:"
for host in node1 node2 n