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
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
| #!/usr/bin/env python3
# monitor_hdfs_balancer.py
import subprocess
import time
import re
from datetime import datetime
def get_hdfs_report():
try:
result = subprocess.run(['hdfs', 'dfsadmin', '-report'],
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
universal_newlines=True, check=True)
return result.stdout
except subprocess.CalledProcessError as e:
print("获取HDFS报告失败: {}".format(e))
return None
def parse_datanode_info(report):
datanodes = []
lines = report.split('\n')
current_node = {}
in_datanode_section = False
for line in lines:
line = line.strip()
if line.startswith('Name:'):
if current_node:
datanodes.append(current_node)
current_node = {'name': line.split(':', 1)[1].strip()}
in_datanode_section = True
elif in_datanode_section and line.startswith('DFS Used%:'):
current_node['used_percent'] = float(line.split(':')[1].strip().replace('%', ''))
elif in_datanode_section and line.startswith('DFS Remaining%:'):
current_node['remaining_percent'] = float(line.split(':')[1].strip().replace('%', ''))
if current_node:
datanodes.append(current_node)
return datanodes
def calculate_balance_metrics(datanodes):
if not datanodes:
return None
used_percents = [node.get('used_percent', 0) for node in datanodes]
avg_used_percent = sum(used_percents) / len(used_percents)
# 计算标准差
variance = sum((x - avg_used_percent) ** 2 for x in used_percents) / len(used_percents)
std_dev = variance ** 0.5
# 找出最高和最低使用率节点
max_used_node = max(datanodes, key=lambda x: x.get('used_percent', 0))
min_used_node = min(datanodes, key=lambda x: x.get('used_percent', 0))
return {
'avg_used_percent': avg_used_percent,
'std_dev': std_dev,
'max_used_node': max_used_node,
'min_used_node': min_used_node,
'datanodes': datanodes
}
def monitor_balancer():
print("HDFS均衡监控开始...")
print("=" * 80)
start_time = datetime.now()
try:
while True:
current_time = datetime.now()
elapsed = current_time - start_time
report = get_hdfs_report()
if not report:
print("[{}] 无法获取HDFS报告".format(current_time.strftime('%H:%M:%S')))
time.sleep(30)
continue
datanodes = parse_datanode_info(report)
if not datanodes:
print("[{}] 无法解析datanode信息".format(current_time.strftime('%H:%M:%S')))
time.sleep(30)
continue
metrics = calculate_balance_metrics(datanodes)
if not metrics:
continue
# 显示当前状态
print("\n[{}] 运行时间: {}".format(current_time.strftime('%H:%M:%S'), elapsed))
print("平均使用率: {:.2f}%".format(metrics['avg_used_percent']))
print("均衡度(标准差): {:.2f}%".format(metrics['std_dev']))
print("最高使用率节点: {} ({:.2f}%)".format(
metrics['max_used_node'].get('name', 'N/A'),
metrics['max_used_node'].get('used_percent', 0)))
print("最低使用率节点: {} ({:.2f}%)".format(
metrics['min_used_node'].get('name', 'N/A'),
metrics['min_used_node'].get('used_percent', 0)))
print("\n各节点使用率:")
for node in sorted(metrics['datanodes'], key=lambda x: x.get('used_percent', 0), reverse=True):
name = node.get('name', 'N/A')
used_pct = node.get('used_percent', 0)
remaining_pct = node.get('remaining_percent', 0)
print(" {}: {:.2f}% (剩余: {:.2f}%)".format(name, used_pct, remaining_pct))
# 检查是否达到均衡
if metrics['std_dev'] < 5.0:
print("\n🎉 均衡完成! 标准差: {:.2f}%".format(metrics['std_dev']))
break
print("=" * 80)
time.sleep(60)
except KeyboardInterrupt:
print("\n\n监控已停止")
except Exception as e:
print("\n监控出错: {}".format(e))
if __name__ == "__main__":
monitor_balancer()
|