From 948305f889420b9b502a9777879227178171cf7b Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Thu, 24 Feb 2022 13:11:33 +0530 Subject: [PATCH 1/7] HDFS-16481. Provide support to set Http and Rpc ports in MiniJournalCluster --- .../hdfs/qjournal/MiniJournalCluster.java | 31 ++++++++- .../hdfs/qjournal/TestMiniJournalCluster.java | 67 +++++++++++++++++++ 2 files changed, 96 insertions(+), 2 deletions(-) diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java index d0bbd44f1afbb1..3836500b18a08c 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java @@ -52,6 +52,8 @@ public static class Builder { private int numJournalNodes = 3; private boolean format = true; private final Configuration conf; + private int[] httpPorts = null; + private int[] rpcPorts = null; static { DefaultMetricsSystem.setMiniClusterMode(true); @@ -76,6 +78,16 @@ public Builder format(boolean f) { return this; } + public Builder setHttpPorts(int... httpPorts) { + this.httpPorts = httpPorts; + return this; + } + + public Builder setRpcPorts(int... rpcPorts) { + this.rpcPorts = rpcPorts; + return this; + } + public MiniJournalCluster build() throws IOException { return new MiniJournalCluster(this); } @@ -99,6 +111,19 @@ private JNInfo(JournalNode node) { private final JNInfo[] nodes; private MiniJournalCluster(Builder b) throws IOException { + + if (b.httpPorts != null && b.httpPorts.length != b.numJournalNodes) { + throw new IllegalArgumentException( + "Num of http ports (" + b.httpPorts.length + ") should match num of JournalNodes (" + + b.numJournalNodes + ")"); + } + + if (b.rpcPorts != null && b.rpcPorts.length != b.numJournalNodes) { + throw new IllegalArgumentException( + "Num of rpc ports (" + b.rpcPorts.length + ") should match num of JournalNodes (" + + b.numJournalNodes + ")"); + } + LOG.info("Starting MiniJournalCluster with " + b.numJournalNodes + " journal nodes"); @@ -173,8 +198,10 @@ private Configuration createConfForNode(Builder b, int idx) { Configuration conf = new Configuration(b.conf); File logDir = getStorageDir(idx); conf.set(DFSConfigKeys.DFS_JOURNALNODE_EDITS_DIR_KEY, logDir.toString()); - conf.set(DFSConfigKeys.DFS_JOURNALNODE_RPC_ADDRESS_KEY, "localhost:0"); - conf.set(DFSConfigKeys.DFS_JOURNALNODE_HTTP_ADDRESS_KEY, "localhost:0"); + int httpPort = b.httpPorts != null ? b.httpPorts[idx] : 0; + int rpcPort = b.rpcPorts != null ? b.rpcPorts[idx] : 0; + conf.set(DFSConfigKeys.DFS_JOURNALNODE_RPC_ADDRESS_KEY, "localhost:" + rpcPort); + conf.set(DFSConfigKeys.DFS_JOURNALNODE_HTTP_ADDRESS_KEY, "localhost:" + httpPort); return conf; } diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java index cace7c92891ab6..61e798657bc55c 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java @@ -31,6 +31,7 @@ public class TestMiniJournalCluster { + @Test public void testStartStop() throws IOException { Configuration conf = new Configuration(); @@ -52,4 +53,70 @@ public void testStartStop() throws IOException { c.shutdown(); } } + + @Test + public void testStartStopWithPorts() throws IOException { + Configuration conf = new Configuration(); + + try { + new MiniJournalCluster.Builder(conf).setHttpPorts(8481).build(); + fail("Should not reach here"); + } catch (IllegalArgumentException e) { + assertEquals("Num of http ports (1) should match num of JournalNodes (3)", e.getMessage()); + } + + try { + new MiniJournalCluster.Builder(conf).setRpcPorts(8481, 8482) + .build(); + fail("Should not reach here"); + } catch (IllegalArgumentException e) { + assertEquals("Num of rpc ports (2) should match num of JournalNodes (3)", e.getMessage()); + } + + try { + new MiniJournalCluster.Builder(conf).setHttpPorts(800, 9000, 10000).setRpcPorts(8481) + .build(); + fail("Should not reach here"); + } catch (IllegalArgumentException e) { + assertEquals("Num of rpc ports (1) should match num of JournalNodes (3)", e.getMessage()); + } + + try { + new MiniJournalCluster.Builder(conf).setHttpPorts(800, 9000, 1000, 2000) + .setRpcPorts(8481, 8482, 8483) + .build(); + fail("Should not reach here"); + } catch (IllegalArgumentException e) { + assertEquals("Num of http ports (4) should match num of JournalNodes (3)", e.getMessage()); + } + + MiniJournalCluster miniJournalCluster = + new MiniJournalCluster.Builder(conf).setHttpPorts(8481, 8482, 8483) + .setRpcPorts(8491, 8492, 8493).build(); + try { + miniJournalCluster.waitActive(); + URI uri = miniJournalCluster.getQuorumJournalURI("myjournal"); + String[] addrs = uri.getAuthority().split(";"); + assertEquals(3, addrs.length); + + assertEquals(8481, miniJournalCluster.getJournalNode(0).getHttpAddress().getPort()); + assertEquals(8482, miniJournalCluster.getJournalNode(1).getHttpAddress().getPort()); + assertEquals(8483, miniJournalCluster.getJournalNode(2).getHttpAddress().getPort()); + + assertEquals(8491, + miniJournalCluster.getJournalNode(0).getRpcServer().getAddress().getPort()); + assertEquals(8492, + miniJournalCluster.getJournalNode(1).getRpcServer().getAddress().getPort()); + assertEquals(8493, + miniJournalCluster.getJournalNode(2).getRpcServer().getAddress().getPort()); + + JournalNode node = miniJournalCluster.getJournalNode(0); + String dir = node.getConf().get(DFSConfigKeys.DFS_JOURNALNODE_EDITS_DIR_KEY); + assertEquals(new File(MiniDFSCluster.getBaseDirectory() + "journalnode-0").getAbsolutePath(), + dir); + } finally { + miniJournalCluster.shutdown(); + } + } + } From 5833956b6f9d4c3116897d027efe6d4cabdbd1c9 Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Thu, 24 Feb 2022 22:16:10 +0530 Subject: [PATCH 2/7] fix checkstyle --- .../apache/hadoop/hdfs/qjournal/MiniJournalCluster.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java index 3836500b18a08c..d858de1bd92d36 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java @@ -78,13 +78,13 @@ public Builder format(boolean f) { return this; } - public Builder setHttpPorts(int... httpPorts) { - this.httpPorts = httpPorts; + public Builder setHttpPorts(int... ports) { + this.httpPorts = ports; return this; } - public Builder setRpcPorts(int... rpcPorts) { - this.rpcPorts = rpcPorts; + public Builder setRpcPorts(int... ports) { + this.rpcPorts = ports; return this; } From e70c3115d9500a1f9e1cf433050f4c47445284eb Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Tue, 1 Mar 2022 18:25:56 +0530 Subject: [PATCH 3/7] addressing reviews from @ayushtkn and @tomscut --- .../hdfs/qjournal/MiniJournalCluster.java | 10 +- .../hdfs/qjournal/TestMiniJournalCluster.java | 103 +++++++++++------- 2 files changed, 71 insertions(+), 42 deletions(-) diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java index d858de1bd92d36..a055a2a3db86a3 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java @@ -20,6 +20,7 @@ import static org.apache.hadoop.hdfs.qjournal.QJMTestUtil.FAKE_NSINFO; import static org.junit.Assert.fail; +import java.io.Closeable; import java.io.File; import java.io.IOException; import java.net.InetSocketAddress; @@ -45,7 +46,8 @@ import org.apache.hadoop.thirdparty.com.google.common.base.Joiner; import org.apache.hadoop.test.GenericTestUtils; -public class MiniJournalCluster { +public class MiniJournalCluster implements Closeable { + public static final String CLUSTER_WAITACTIVE_URI = "waitactive"; public static class Builder { private String baseDir; @@ -301,4 +303,10 @@ public void setNamenodeSharedEditsConf(String jid) { .DFS_NAMENODE_SHARED_EDITS_DIR_KEY, quorumJournalURI.toString()); } } + + @Override + public void close() throws IOException { + this.shutdown(); + } + } diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java index 61e798657bc55c..f114feddb64e56 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java @@ -27,11 +27,16 @@ import org.apache.hadoop.hdfs.DFSConfigKeys; import org.apache.hadoop.hdfs.MiniDFSCluster; import org.apache.hadoop.hdfs.qjournal.server.JournalNode; -import org.junit.Test; +import org.apache.hadoop.test.LambdaTestUtils; +import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; public class TestMiniJournalCluster { + private static final Logger LOG = LoggerFactory.getLogger(TestMiniJournalCluster.class); + @Test public void testStartStop() throws IOException { Configuration conf = new Configuration(); @@ -55,67 +60,83 @@ public void testStartStop() throws IOException { } @Test - public void testStartStopWithPorts() throws IOException { + public void testStartStopWithPorts() throws Exception { Configuration conf = new Configuration(); - try { - new MiniJournalCluster.Builder(conf).setHttpPorts(8481).build(); - fail("Should not reach here"); - } catch (IllegalArgumentException e) { - assertEquals("Num of http ports (1) should match num of JournalNodes (3)", e.getMessage()); - } + LambdaTestUtils.intercept( + IllegalArgumentException.class, + "Num of http ports (1) should match num of JournalNodes (3)", + "MiniJournalCluster port validation failed", + () -> { + new MiniJournalCluster.Builder(conf).setHttpPorts(8481).build(); + }); + + LambdaTestUtils.intercept( + IllegalArgumentException.class, + "Num of rpc ports (2) should match num of JournalNodes (3)", + "MiniJournalCluster port validation failed", + () -> { + new MiniJournalCluster.Builder(conf).setRpcPorts(8481, 8482).build(); + }); + + LambdaTestUtils.intercept( + IllegalArgumentException.class, + "Num of rpc ports (1) should match num of JournalNodes (3)", + "MiniJournalCluster port validation failed", + () -> { + new MiniJournalCluster.Builder(conf).setHttpPorts(800, 9000, 10000).setRpcPorts(8481) + .build(); + }); + + LambdaTestUtils.intercept( + IllegalArgumentException.class, + "Num of http ports (4) should match num of JournalNodes (3)", + "MiniJournalCluster port validation failed", + () -> { + new MiniJournalCluster.Builder(conf).setHttpPorts(800, 9000, 1000, 2000) + .setRpcPorts(8481, 8482, 8483).build(); + }); + + final int[] httpPorts = new int[3]; + final int[] rpcPorts = new int[3]; + try (MiniJournalCluster miniJournalCluster = new MiniJournalCluster.Builder(conf).build()) { + miniJournalCluster.waitActive(); - try { - new MiniJournalCluster.Builder(conf).setRpcPorts(8481, 8482) - .build(); - fail("Should not reach here"); - } catch (IllegalArgumentException e) { - assertEquals("Num of rpc ports (2) should match num of JournalNodes (3)", e.getMessage()); - } + for (int i = 0; i < 3; i++) { + httpPorts[i] = miniJournalCluster.getJournalNode(i).getHttpAddress().getPort(); + } - try { - new MiniJournalCluster.Builder(conf).setHttpPorts(800, 9000, 10000).setRpcPorts(8481) - .build(); - fail("Should not reach here"); - } catch (IllegalArgumentException e) { - assertEquals("Num of rpc ports (1) should match num of JournalNodes (3)", e.getMessage()); + for (int i = 0; i < 3; i++) { + rpcPorts[i] = miniJournalCluster.getJournalNode(i).getRpcServer().getAddress().getPort(); + } } - try { - new MiniJournalCluster.Builder(conf).setHttpPorts(800, 9000, 1000, 2000) - .setRpcPorts(8481, 8482, 8483) - .build(); - fail("Should not reach here"); - } catch (IllegalArgumentException e) { - assertEquals("Num of http ports (4) should match num of JournalNodes (3)", e.getMessage()); - } + LOG.info("Http ports selected: {}", httpPorts); + LOG.info("Rpc ports selected: {}", rpcPorts); - MiniJournalCluster miniJournalCluster = - new MiniJournalCluster.Builder(conf).setHttpPorts(8481, 8482, 8483) - .setRpcPorts(8491, 8492, 8493).build(); - try { + try (MiniJournalCluster miniJournalCluster = new MiniJournalCluster.Builder(conf) + .setHttpPorts(httpPorts) + .setRpcPorts(rpcPorts).build()) { miniJournalCluster.waitActive(); URI uri = miniJournalCluster.getQuorumJournalURI("myjournal"); String[] addrs = uri.getAuthority().split(";"); assertEquals(3, addrs.length); - assertEquals(8481, miniJournalCluster.getJournalNode(0).getHttpAddress().getPort()); - assertEquals(8482, miniJournalCluster.getJournalNode(1).getHttpAddress().getPort()); - assertEquals(8483, miniJournalCluster.getJournalNode(2).getHttpAddress().getPort()); + assertEquals(httpPorts[0], miniJournalCluster.getJournalNode(0).getHttpAddress().getPort()); + assertEquals(httpPorts[1], miniJournalCluster.getJournalNode(1).getHttpAddress().getPort()); + assertEquals(httpPorts[2], miniJournalCluster.getJournalNode(2).getHttpAddress().getPort()); - assertEquals(8491, + assertEquals(rpcPorts[0], miniJournalCluster.getJournalNode(0).getRpcServer().getAddress().getPort()); - assertEquals(8492, + assertEquals(rpcPorts[1], miniJournalCluster.getJournalNode(1).getRpcServer().getAddress().getPort()); - assertEquals(8493, + assertEquals(rpcPorts[2], miniJournalCluster.getJournalNode(2).getRpcServer().getAddress().getPort()); JournalNode node = miniJournalCluster.getJournalNode(0); String dir = node.getConf().get(DFSConfigKeys.DFS_JOURNALNODE_EDITS_DIR_KEY); assertEquals(new File(MiniDFSCluster.getBaseDirectory() + "journalnode-0").getAbsolutePath(), dir); - } finally { - miniJournalCluster.shutdown(); } } From 24fa464e7aecb792b291c70ab99fb8aaa4d64610 Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Tue, 1 Mar 2022 20:09:30 +0530 Subject: [PATCH 4/7] addressing recent review --- .../hdfs/qjournal/TestMiniJournalCluster.java | 23 ++++++++----------- 1 file changed, 10 insertions(+), 13 deletions(-) diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java index f114feddb64e56..95442aec7a076e 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java @@ -27,6 +27,7 @@ import org.apache.hadoop.hdfs.DFSConfigKeys; import org.apache.hadoop.hdfs.MiniDFSCluster; import org.apache.hadoop.hdfs.qjournal.server.JournalNode; +import org.apache.hadoop.net.NetUtils; import org.apache.hadoop.test.LambdaTestUtils; import org.junit.Test; @@ -97,23 +98,19 @@ public void testStartStopWithPorts() throws Exception { .setRpcPorts(8481, 8482, 8483).build(); }); - final int[] httpPorts = new int[3]; - final int[] rpcPorts = new int[3]; - try (MiniJournalCluster miniJournalCluster = new MiniJournalCluster.Builder(conf).build()) { - miniJournalCluster.waitActive(); - - for (int i = 0; i < 3; i++) { - httpPorts[i] = miniJournalCluster.getJournalNode(i).getHttpAddress().getPort(); - } - - for (int i = 0; i < 3; i++) { - rpcPorts[i] = miniJournalCluster.getJournalNode(i).getRpcServer().getAddress().getPort(); - } - } + final int[] httpPorts = new int[] { NetUtils.getFreeSocketPort(), NetUtils.getFreeSocketPort(), + NetUtils.getFreeSocketPort() }; + final int[] rpcPorts = new int[] { NetUtils.getFreeSocketPort(), NetUtils.getFreeSocketPort(), + NetUtils.getFreeSocketPort() }; LOG.info("Http ports selected: {}", httpPorts); LOG.info("Rpc ports selected: {}", rpcPorts); + for (int i = 0; i < 3; i++) { + assertNotEquals(0, rpcPorts[i]); + assertNotEquals(0, httpPorts[i]); + } + try (MiniJournalCluster miniJournalCluster = new MiniJournalCluster.Builder(conf) .setHttpPorts(httpPorts) .setRpcPorts(rpcPorts).build()) { From bcf8a15d444a92575f341ca6b69858c403667421 Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Wed, 2 Mar 2022 15:58:30 +0530 Subject: [PATCH 5/7] adding new utility to acquire free ports --- .../java/org/apache/hadoop/net/NetUtils.java | 19 +++++++++++++++++++ .../hdfs/qjournal/MiniJournalCluster.java | 2 +- .../hdfs/qjournal/TestMiniJournalCluster.java | 19 +++++++++++++++---- 3 files changed, 35 insertions(+), 5 deletions(-) diff --git a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java index 4b924af03c1966..04337fda46910e 100644 --- a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java +++ b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java @@ -1053,6 +1053,25 @@ public static int getFreeSocketPort() { return port; } + /** + * Return free ports. There is no guarantee they will remain free, so + * ports should be used immediately. The number of free ports returned by + * this method should match argument {@code numOfPorts}. + * + * @param numOfPorts Number of free ports to acquire. + * @return Free ports for binding a local socket. + */ + public static Set getFreeSocketPorts(int numOfPorts) { + final Set freePorts = new HashSet<>(numOfPorts); + for (int i = 0; i < numOfPorts * 5; i++) { + freePorts.add(getFreeSocketPort()); + if (freePorts.size() == numOfPorts) { + return freePorts; + } + } + throw new IllegalStateException(numOfPorts + " free ports could not be acquired."); + } + /** * Return an @{@link InetAddress} to bind to. If bindWildCardAddress is true * than returns null. diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java index a055a2a3db86a3..1c43b39159a99f 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/MiniJournalCluster.java @@ -46,7 +46,7 @@ import org.apache.hadoop.thirdparty.com.google.common.base.Joiner; import org.apache.hadoop.test.GenericTestUtils; -public class MiniJournalCluster implements Closeable { +public final class MiniJournalCluster implements Closeable { public static final String CLUSTER_WAITACTIVE_URI = "waitactive"; public static class Builder { diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java index 95442aec7a076e..daff40d898f5c4 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java @@ -22,6 +22,7 @@ import java.io.File; import java.io.IOException; import java.net.URI; +import java.util.Set; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hdfs.DFSConfigKeys; @@ -98,10 +99,20 @@ public void testStartStopWithPorts() throws Exception { .setRpcPorts(8481, 8482, 8483).build(); }); - final int[] httpPorts = new int[] { NetUtils.getFreeSocketPort(), NetUtils.getFreeSocketPort(), - NetUtils.getFreeSocketPort() }; - final int[] rpcPorts = new int[] { NetUtils.getFreeSocketPort(), NetUtils.getFreeSocketPort(), - NetUtils.getFreeSocketPort() }; + final Set httpAndRpcPorts = NetUtils.getFreeSocketPorts(6); + LOG.info("Free socket ports: {}", httpAndRpcPorts); + + final int[] httpPorts = new int[3]; + final int[] rpcPorts = new int[3]; + int httpPortIdx = 0; + int rpcPortIdx = 0; + for (Integer httpAndRpcPort : httpAndRpcPorts) { + if (httpPortIdx < 3) { + httpPorts[httpPortIdx++] = httpAndRpcPort; + } else { + rpcPorts[rpcPortIdx++] = httpAndRpcPort; + } + } LOG.info("Http ports selected: {}", httpPorts); LOG.info("Rpc ports selected: {}", rpcPorts); From 68e7e2eb60f236b7b59a383ab2b8da06e724648d Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Thu, 3 Mar 2022 11:37:44 +0530 Subject: [PATCH 6/7] exclude 0 in the free ports --- .../src/main/java/org/apache/hadoop/net/NetUtils.java | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java index 04337fda46910e..fead87d7907d76 100644 --- a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java +++ b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/net/NetUtils.java @@ -1056,15 +1056,22 @@ public static int getFreeSocketPort() { /** * Return free ports. There is no guarantee they will remain free, so * ports should be used immediately. The number of free ports returned by - * this method should match argument {@code numOfPorts}. + * this method should match argument {@code numOfPorts}. Num of ports + * provided in the argument should not exceed 25. * * @param numOfPorts Number of free ports to acquire. * @return Free ports for binding a local socket. */ public static Set getFreeSocketPorts(int numOfPorts) { + Preconditions.checkArgument(numOfPorts > 0 && numOfPorts <= 25, + "Valid range for num of ports is between 0 and 26"); final Set freePorts = new HashSet<>(numOfPorts); for (int i = 0; i < numOfPorts * 5; i++) { - freePorts.add(getFreeSocketPort()); + int port = getFreeSocketPort(); + if (port == 0) { + continue; + } + freePorts.add(port); if (freePorts.size() == numOfPorts) { return freePorts; } From 84e430e4cc436387f66b860d49a6d6d14b7378e3 Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Thu, 3 Mar 2022 20:44:28 +0530 Subject: [PATCH 7/7] assert all acquire ports --- .../hadoop/hdfs/qjournal/TestMiniJournalCluster.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java index daff40d898f5c4..ccbbc94c99ede1 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/TestMiniJournalCluster.java @@ -102,6 +102,11 @@ public void testStartStopWithPorts() throws Exception { final Set httpAndRpcPorts = NetUtils.getFreeSocketPorts(6); LOG.info("Free socket ports: {}", httpAndRpcPorts); + for (Integer httpAndRpcPort : httpAndRpcPorts) { + assertNotEquals("None of the acquired socket port should not be zero", 0, + httpAndRpcPort.intValue()); + } + final int[] httpPorts = new int[3]; final int[] rpcPorts = new int[3]; int httpPortIdx = 0; @@ -117,11 +122,6 @@ public void testStartStopWithPorts() throws Exception { LOG.info("Http ports selected: {}", httpPorts); LOG.info("Rpc ports selected: {}", rpcPorts); - for (int i = 0; i < 3; i++) { - assertNotEquals(0, rpcPorts[i]); - assertNotEquals(0, httpPorts[i]); - } - try (MiniJournalCluster miniJournalCluster = new MiniJournalCluster.Builder(conf) .setHttpPorts(httpPorts) .setRpcPorts(rpcPorts).build()) {