diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/DFSConfigKeys.java b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/DFSConfigKeys.java
index f92a2ad56581b9..8cb90bd47e2fd2 100755
--- a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/DFSConfigKeys.java
+++ b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/DFSConfigKeys.java
@@ -293,6 +293,9 @@ public class DFSConfigKeys extends CommonConfigurationKeys {
public static final boolean
DFS_NAMENODE_REDUNDANCY_CONSIDERLOADBYVOLUME_DEFAULT
= false;
+ public static final String DFS_NAMENODE_REDUNDANCY_CONSIDERLOAD_MINLOAD_KEY =
+ "dfs.namenode.redundancy.considerLoad.minload";
+ public static final int DFS_NAMENODE_REDUNDANCY_CONSIDERLOAD_MINLOAD_DEFAULT = 16;
public static final String DFS_NAMENODE_REDUNDANCY_INTERVAL_SECONDS_KEY =
HdfsClientConfigKeys.DeprecatedKeys.DFS_NAMENODE_REDUNDANCY_INTERVAL_SECONDS_KEY;
public static final int DFS_NAMENODE_REDUNDANCY_INTERVAL_SECONDS_DEFAULT = 3;
diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/blockmanagement/BlockPlacementPolicyDefault.java b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/blockmanagement/BlockPlacementPolicyDefault.java
index 8020d7c45b37ad..348aac33bfcaf2 100644
--- a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/blockmanagement/BlockPlacementPolicyDefault.java
+++ b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/blockmanagement/BlockPlacementPolicyDefault.java
@@ -105,6 +105,7 @@ private String getText() {
private boolean considerLoadByStorageType;
protected double considerLoadFactor;
private boolean considerLoadByVolume = false;
+ protected int considerLoadMinLoad;
private boolean preferLocalNode;
private boolean dataNodePeerStatsEnabled;
private volatile boolean excludeSlowNodesEnabled;
@@ -140,6 +141,9 @@ public void initialize(Configuration conf, FSClusterStats stats,
DFSConfigKeys.DFS_NAMENODE_REDUNDANCY_CONSIDERLOADBYVOLUME_KEY,
DFSConfigKeys.DFS_NAMENODE_REDUNDANCY_CONSIDERLOADBYVOLUME_DEFAULT
);
+ this.considerLoadMinLoad = conf.getInt(
+ DFSConfigKeys.DFS_NAMENODE_REDUNDANCY_CONSIDERLOAD_MINLOAD_KEY,
+ DFSConfigKeys.DFS_NAMENODE_REDUNDANCY_CONSIDERLOAD_MINLOAD_DEFAULT);
this.stats = stats;
this.clusterMap = clusterMap;
this.host2datanodeMap = host2datanodeMap;
@@ -1014,7 +1018,7 @@ boolean excludeNodeByLoad(DatanodeDescriptor node){
final double maxLoad = considerLoadFactor * inServiceXceiverCount;
final int nodeLoad = node.getXceiverCount();
- if ((nodeLoad > maxLoad) && (maxLoad > 0)) {
+ if ((nodeLoad > considerLoadMinLoad) &&(nodeLoad > maxLoad) && (maxLoad > 0)) {
logNodeIsNotChosen(node, NodeNotChosenReason.NODE_TOO_BUSY,
"(load: " + nodeLoad + " > " + maxLoad + ")");
return true;
diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/main/resources/hdfs-default.xml b/hadoop-hdfs-project/hadoop-hdfs/src/main/resources/hdfs-default.xml
index e6dc8c5ba1ac42..13f2eb9f371cf1 100755
--- a/hadoop-hdfs-project/hadoop-hdfs/src/main/resources/hdfs-default.xml
+++ b/hadoop-hdfs-project/hadoop-hdfs/src/main/resources/hdfs-default.xml
@@ -334,6 +334,14 @@
+
+ dfs.namenode.redundancy.considerLoad.minload
+ 16
+ The minimum load which a node's load must exceed
+ before being rejected for writes, only if considerLoad is true.
+
+
+
dfs.namenode.redundancy.considerLoadByVolume
false
diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/server/blockmanagement/TestReplicationPolicy.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/server/blockmanagement/TestReplicationPolicy.java
index b99e060ee387ca..868cc0d8c4a348 100644
--- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/server/blockmanagement/TestReplicationPolicy.java
+++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/server/blockmanagement/TestReplicationPolicy.java
@@ -1662,6 +1662,7 @@ public void testMaxLoad() {
when(node.getXceiverCount()).thenReturn(1);
final Configuration conf = new Configuration();
+ conf.setInt(DFSConfigKeys.DFS_NAMENODE_REDUNDANCY_CONSIDERLOAD_MINLOAD_KEY, 1);
final Class extends BlockPlacementPolicy> replicatorClass = conf
.getClass(DFSConfigKeys.DFS_BLOCK_REPLICATOR_CLASSNAME_KEY,
DFSConfigKeys.DFS_BLOCK_REPLICATOR_CLASSNAME_DEFAULT,
@@ -1708,6 +1709,32 @@ public void testMaxLoad() {
assertFalse(bppd.excludeNodeByLoad(node));
}
+ @Test
+ public void testMinLoad() {
+ FSClusterStats statistics = mock(FSClusterStats.class);
+ DatanodeDescriptor node = mock(DatanodeDescriptor.class);
+
+ when(statistics.getInServiceXceiverAverage()).thenReturn(5D);
+ when(node.getXceiverCount()).thenReturn(12);
+
+ final Configuration conf = new Configuration();
+ conf.setInt(DFSConfigKeys.DFS_NAMENODE_REDUNDANCY_CONSIDERLOAD_MINLOAD_KEY, 16);
+ final Class extends BlockPlacementPolicy> replicatorClass = conf
+ .getClass(DFSConfigKeys.DFS_BLOCK_REPLICATOR_CLASSNAME_KEY,
+ DFSConfigKeys.DFS_BLOCK_REPLICATOR_CLASSNAME_DEFAULT,
+ BlockPlacementPolicy.class);
+ BlockPlacementPolicy bpp = ReflectionUtils.
+ newInstance(replicatorClass, conf);
+ assertTrue(bpp instanceof BlockPlacementPolicyDefault);
+
+ BlockPlacementPolicyDefault bppd = (BlockPlacementPolicyDefault) bpp;
+ bppd.initialize(conf, statistics, null, null);
+ assertFalse(bppd.excludeNodeByLoad(node));
+
+ when(node.getXceiverCount()).thenReturn(17);
+ assertTrue(bppd.excludeNodeByLoad(node));
+ }
+
@Test
public void testChosenFailureForStorageType() {
final LogVerificationAppender appender = new LogVerificationAppender();