Recommended HDFS configurations in E-MapReduce (EMR) to improve performance and stability.
Background
Apply the following HDFS optimizations as needed:
Control the number of small files
-
Context: The HDFS NameNode loads all file metadata into memory. In a cluster with a fixed disk capacity, an excessive number of small files can cause a memory bottleneck on the NameNode.
-
Recommendation: Limit the number of small files. Merge existing small files into larger ones where possible.
Configure the file limit for HDFS directories
-
Context: When a cluster is running, different components, such as Spark and YARN, or clients might continuously write files to the same HDFS directory. However, HDFS has a limit on the number of files a single directory can contain. To prevent job failures, plan your storage so you do not exceed this limit.
-
Recommendation: In the EMR console, go to the Configure tab for the HDFS service. Click the hdfs-site tab, and then click Add Configuration Item. Add the dfs.namenode.fs-limits.max-directory-items parameter to set the maximum number of items per directory, and then save the configuration. For more information about how to add a parameter, see Manage configuration items.
NotePlan your data storage by organizing files into categories, such as by date or business type. Avoid storing too many items directly in a single directory. Keep the number of items in a single directory to approximately one million.
Configure the tolerance for damaged volumes
-
Context: By default, a DataNode with multiple storage volumes stops providing service if one volume becomes damaged.
-
Recommendation: In the EMR console, go to the Configure tab for the HDFS service. In the Configuration Search area, search for and modify the dfs.datanode.failed.volumes.tolerated parameter. This parameter sets the number of volume failures a DataNode can tolerate. The DataNode remains operational as long as the number of failed volumes does not exceed this limit.
Parameter
Description
Default
dfs.datanode.failed.volumes.tolerated
The maximum number of volume failures a DataNode can tolerate before it stops providing service.
0
NoteIf a DataNode has a damaged disk and no other nodes are available to take over, you can temporarily increase this value to bring the DataNode back online.
Use the balancer for capacity balancing
-
Context: Disk usage among DataNodes in an HDFS cluster can become imbalanced, especially after adding new nodes. This imbalance can overload some DataNodes.
-
Recommendation: Run the HDFS balancer to redistribute data evenly across DataNodes.
NoteRunning the balancer consumes network bandwidth on the DataNodes. To avoid disrupting active workloads, run the balancer during off-peak hours.
-
Log on to any node in the cluster.
-
Optional: Run the following command to change the maximum bandwidth for the balancer.
hdfs dfsadmin -setBalancerBandwidth <bandwidth in bytes per second>NoteIn the command,
<bandwidth in bytes per second>is the maximum bandwidth. For example, to set the bandwidth limit to 200 MB/s, the value is 200 * 1024 * 1024, which is 209715200 bytes. The complete command ishdfs dfsadmin -setBalancerBandwidth 209715200. If the cluster load is high, reduce the bandwidth limit to relieve network pressure, such as to 20971520 (20 MB/s). If the cluster is idle, increase the limit to accelerate data balancing, such as to 1073741824 (1 GB/s). -
Run the following commands to switch to the hdfs user and start the balancer.
su hdfs /opt/apps/HDFS/hdfs-current/sbin/start-balancer.sh -threshold 10 -
Run the following command to view the balancer log.
tail -f /var/log/emr/hadoop-hdfs/hadoop-hdfs-balancer-master-1-1.c-xxx.logNoteIn the command,
hadoop-hdfs-balancer-master-1-1.c-xxxis a placeholder for the log file name generated by the balancer.
-
The word Successfully in the output indicates that the operation succeeded.
Balance DataNodes with different disk specifications
-
Context: When a cluster contains DataNode node groups with different disk specifications, such as a different disk size or a different number of disks per node, the balancer equalizes disk usage as a percentage. As a result, the usage percentages converge across node groups, but the absolute HDFS available disk space still differs between them.
-
Recommendation: Set a different dfs.datanode.du.reserved value for each node group so that all node groups provide the same HDFS available disk space. This parameter specifies the space reserved on each disk for non-HDFS use, which is excluded from the HDFS available disk space.
Parameter
Description
Default
dfs.datanode.du.reserved
The space reserved on each disk for non-HDFS use. This space is excluded from the HDFS available disk space. Unit: bytes.
1073741824 (1 GB)
NoteThis parameter supports node-group-level configuration, so you can set a different value for each node group that has a different disk specification.
Use the following formula to calculate the HDFS available disk space per DataNode:
HDFS available disk space per DataNode = (disk size - dfs.datanode.du.reserved) x number of disks
For example, node group A has 4 disks of 1200 GB each, with dfs.datanode.du.reserved set to 380 GB. Node group B has 8 disks of 7300 GB each, with dfs.datanode.du.reserved set to 6890 GB. Both node groups then provide 3280 GB of HDFS available disk space per DataNode:
-
Node group A: (1200 - 380) x 4 = 3280 GB
-
Node group B: (7300 - 6890) x 8 = 3280 GB
To configure dfs.datanode.du.reserved for each node group:
-
Log on to the EMR on ECS console, and go to the Configure tab of the HDFS service for the target cluster.
-
In the Configuration Search area, search for dfs.datanode.du.reserved, switch the configuration scope to the target node group, and enter the value calculated with the formula above.
-
Repeat the previous step for every node group that has a different disk specification, and then save the configuration.
-
Restart the HDFS service for the change to take effect.
-
On the HDFS web UI, check whether the capacity of each DataNode is now consistent. If the capacities are consistent but data is still unevenly distributed, run
hdfs balanceras described in the previous section.
NoteEnter the value in bytes, or in the unit required by the console. Convert the GB values in the example to the corresponding number of bytes. Make sure that (disk size - reserved value) is positive and that enough space is left for the operating system, logs, and other non-HDFS use.
-