The HDFS balancer analyzes block distribution and redistributes data among DataNodes to maintain even storage utilization across your cluster.
Background
HDFS uses a master-slave architecture. The NameNode manages the file system's metadata, such as file names and the location of data blocks, while multiple DataNodes store the actual data blocks. This architecture enables data redundancy and improves the system's fault tolerance.
Over time, operations like adding, deleting, and modifying files can lead to an uneven distribution of data across DataNodes. Some nodes may approach their storage capacity, while others have ample free space. This imbalance affects storage efficiency and increases the risk of data loss, as overutilized nodes are more vulnerable to hardware failure.
To address this issue, HDFS provides the HDFS balancer, a command-line utility that moves data blocks across DataNodes to reduce storage imbalance and improve resource utilization.
View DataNode capacity and usage
Check DataNode capacity and usage to identify storage imbalances and determine whether rebalancing is needed.
-
Connect to the master node of the target cluster. For more information, see Connect to a cluster.
-
Run the following command to view the capacity and usage of each DataNode.
hdfs dfsadmin -report
The output shows each DataNode's total capacity, used space, usage percentage, and remaining space.
Start the HDFS balancer when storage usage differences between DataNodes exceed the configured threshold (10% by default).
Start the HDFS balancer
Method 1: Use the hdfs balancer command
The syntax for the hdfs balancer command is as follows:
hdfs balancer
[-threshold <threshold>]
[-policy <policy>]
[-exclude [-f <hosts-file> | <comma-separated list of hosts>]]
[-include [-f <hosts-file> | <comma-separated list of hosts>]]
[-source [-f <hosts-file> | <comma-separated list of hosts>]]
[-blockpools <comma-separated list of blockpool ids>]
[-idleiterations <idleiterations>]
The following table describes the main parameters for the balancer.
|
Parameter |
Description |
|
threshold |
The allowed percentage of deviation from the average disk usage across the cluster. The default value is 10%, meaning the usage of any DataNode should be within 10% (plus or minus) of the cluster average. When the overall cluster usage is high, you should lower this value to prevent the balancer from failing to find movable blocks. When you add many new nodes to the cluster, you can increase this value to accelerate data movement from high-usage nodes to low-usage nodes. |
|
policy |
The balancing policy. The following policies are supported:
|
|
exclude |
Excludes specified DataNodes from the balancing operation. |
|
include |
Includes only the specified DataNodes in the balancing operation. |
|
source |
Specifies which DataNodes can be sources for data movement. |
|
blockpools |
Runs the balancer only on the specified block pools. |
|
idleiterations |
The maximum number of idle iterations before the balancer exits. This overrides the default of 5. |
Method 2: Use the start-balancer.sh tool
The start-balancer.sh script is a wrapper that calls the hdfs daemon start balancer command. To use the script, follow these steps:
-
Connect to any node in the target cluster. For more information, see Connect to a cluster.
-
Optional: Run the following command to set the balancer's maximum bandwidth.
hdfs dfsadmin -setBalancerBandwidth <bandwidth in bytes per second>NoteThe
<bandwidth in bytes per second>parameter specifies the maximum bandwidth. For example, to set the bandwidth to 200 MB/s, the value is 209,715,200 (200 * 1024 * 1024 bytes). The full command ishdfs dfsadmin -setBalancerBandwidth 209715200. To optimize network usage and protect critical services during high-load periods, we recommend reducing the balancing bandwidth, for example to 20,971,520 (20 MB/s). During off-peak hours, you can increase the bandwidth, for example to 1,073,741,824 (1 GB/s), to speed up the balancing process. -
Run the following command to switch to the
hdfsuser and execute the balancer script.-
DataLake cluster
su hdfs /opt/apps/HDFS/hdfs-current/sbin/start-balancer.sh -threshold 5 -
Hadoop cluster
su hdfs /usr/lib/hadoop-current/sbin/start-balancer.sh -threshold 5NoteThe
-threshold 5parameter sets the balancing threshold. A threshold of 5% means the balancer considers a DataNode balanced if its storage usage is within 5% of the cluster's average usage. The balancer will not move data blocks to or from that node. You can adjust this value to achieve the desired balance for your environment.
-
-
Run the following command to monitor the balancer's progress.
-
DataLake cluster
tail -f /var/log/emr/hadoop-hdfs/hadoop-hdfs-balancer-master-1-1.c-xxx.log -
Hadoop cluster
tail -f /var/log/hadoop-hdfs/hadoop-hdfs-balancer-emr-header-1.cluster-xxx.logNoteThe log file names
hadoop-hdfs-balancer-master-1-1.c-xxx.logandhadoop-hdfs-balancer-emr-header-xx.cluster-xxx.logare examples. Use the actual log file names from your cluster.
If the output contains the word
Successfully, the operation succeeded. -
Tune balancer parameters
The balancer consumes system resources, so run it during off-peak hours. By default, no parameter adjustment is needed. If tuning is necessary, navigate to the HDFS service page in the E-MapReduce console, select , and adjust the following configurations.
-
Client configurations
Parameter
Description
dfs.balancer.dispatcherThreads
Before moving data blocks, the balancer retrieves a list of blocks in each iteration and dispatches them to mover threads.
NoteThe dispatcherThreads parameter defines the number of threads for this dispatching task. The default is 200.
dfs.balancer.rpc.per.sec
The number of RPCs sent per second. The default is 20.
The dispatcher threads make numerous
getBlocksRPC calls to the NameNode. To prevent excessive load on the NameNode, you can control the rate of these RPCs.For example, on a high-load cluster, you can reduce this value to 10 or 5. This has a minimal impact on the overall progress of the data movement.
dfs.balancer.getBlocks.size
In each iteration, the balancer retrieves a list of blocks for the mover threads. This parameter sets the total size of blocks in that list, with a default of 2 GB. Because the
getBlocksprocess locks the NameNode, you can adjust this value based on the NameNode's current load.dfs.balancer.moverThreads
The number of threads that move data blocks. Each thread handles one block move.
The default value is 1000.
-
DataNode configurations
Parameter
Description
dfs.datanode.balance.bandwidthPerSec
Specifies the bandwidth that a DataNode can use for balancing. The recommended setting is 100 MB/s. You can also adjust this value dynamically with the dfsadmin -setBalancerBandwidth command without restarting the DataNode.
For example, increase the bandwidth during low-load periods and decrease it during high-load periods.
dfs.datanode.balance.max.concurrent.moves
The maximum number of blocks a DataNode can move concurrently. The default is 5.
Typically, this value should correspond to the number of disks. We recommend setting an upper limit of
4 *on the DataNode side and then tuning it using the balancer.For example, if a DataNode has 28 disks, you can set the value to 28 on the balancer side and to
28 * 4on the DataNode side. Adjust the values based on the cluster load. Increase the number of concurrent moves when the load is low, and decrease it when the load is high.
FAQ
Q: Why is the usage difference between nodes still around 20% after balancing, even when the threshold is set to 10%?
A: The threshold parameter ensures that the usage of each DataNode does not deviate from the cluster's average usage by more than the specified percentage. Therefore, after balancing, the usage gap between the most- and least-utilized nodes can be up to twice the threshold. For example, one node might be 10% above the average, while another is 10% below, resulting in a 20% gap between them. To reduce this gap, try setting the threshold to a smaller value, such as 5%.