This is an automated email from the ASF dual-hosted git repository.
peterxcli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 0f6295ef3ab HDDS-15085. Add DN and Cluster Readiness (#10757)
0f6295ef3ab is described below
commit 0f6295ef3abc613f6bd3d5093a5cb10cde11f3b3
Author: Chun-Hung Tseng <[email protected]>
AuthorDate: Mon Jul 20 19:02:29 2026 +0200
HDDS-15085. Add DN and Cluster Readiness (#10757)
Co-authored-by: Bolin Lin <[email protected]>
Co-authored-by: Peter Lee <[email protected]>
---
.../ozone/local/TestLocalOzoneClusterRuntime.java | 32 ++-
hadoop-ozone/tools/pom.xml | 4 +
.../hadoop/ozone/local/LocalOzoneCluster.java | 228 +++++++++++++++++++--
.../hadoop/ozone/local/TestLocalOzoneCluster.java | 109 ++++++++++
4 files changed, 350 insertions(+), 23 deletions(-)
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
index 2d2d762896a..a4a10f0491b 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
@@ -17,6 +17,7 @@
package org.apache.hadoop.ozone.local;
+import static java.nio.charset.StandardCharsets.UTF_8;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -37,25 +38,28 @@
import org.junit.jupiter.api.io.TempDir;
/**
- * Integration tests for the SCM and OM portion of {@link LocalOzoneCluster}.
+ * Integration tests for {@link LocalOzoneCluster}.
*/
class TestLocalOzoneClusterRuntime {
+ private static final String KEY_CONTENT = "local ozone key content";
+
@TempDir
private Path tempDir;
@Test
- void scmAndOmStartAndReuseExistingMetadata() throws Exception {
+ void clusterStartsAndReusesExistingData() throws Exception {
String volumeName = uniqueName("vol");
String bucketName = uniqueName("bucket");
+ String keyName = uniqueName("key");
Path dataDir = tempDir.resolve("local-ozone-runtime");
LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
.setS3gEnabled(false)
.setStartupTimeout(Duration.ofMinutes(2))
.build();
- startRuntimeAndCreateBucket(config, volumeName, bucketName);
- restartRuntimeAndVerifyBucket(config, volumeName, bucketName);
+ startRuntimeAndCreateKey(config, volumeName, bucketName, keyName);
+ restartRuntimeAndVerifyKey(config, volumeName, bucketName, keyName);
}
@Test
@@ -83,34 +87,44 @@ void formatNeverRejectsUninitializedScmOmStorage() throws
Exception {
error.getMessage());
}
- private void startRuntimeAndCreateBucket(LocalOzoneClusterConfig config,
- String volumeName, String bucketName) throws Exception {
+ private void startRuntimeAndCreateKey(LocalOzoneClusterConfig config,
+ String volumeName, String bucketName, String keyName) throws Exception {
try (LocalOzoneCluster cluster = new LocalOzoneCluster(config, new
OzoneConfiguration())) {
OzoneConfiguration clientConf =
cluster.prepareConfiguration().getConfiguration();
cluster.start();
+ assertEquals(config.getDatanodes(), cluster.getDatanodeCount());
assertServicePortsReachable(cluster);
try (OzoneClient client = OzoneClientFactory.getRpcClient(clientConf)) {
- TestDataUtil.createVolumeAndBucket(client, volumeName, bucketName);
+ OzoneBucket bucket =
+ TestDataUtil.createVolumeAndBucket(client, volumeName, bucketName);
+ // Writing and reading back a key proves the datanodes registered and
+ // SCM left safe mode, so the cluster is actually usable.
+ TestDataUtil.createKey(bucket, keyName, KEY_CONTENT.getBytes(UTF_8));
+ assertEquals(KEY_CONTENT, TestDataUtil.getKey(bucket, keyName));
}
}
}
- private void restartRuntimeAndVerifyBucket(LocalOzoneClusterConfig config,
- String volumeName, String bucketName) throws Exception {
+ private void restartRuntimeAndVerifyKey(LocalOzoneClusterConfig config,
+ String volumeName, String bucketName, String keyName) throws Exception {
try (LocalOzoneCluster cluster = new LocalOzoneCluster(config, new
OzoneConfiguration())) {
OzoneConfiguration clientConf =
cluster.prepareConfiguration().getConfiguration();
cluster.start();
+ assertEquals(config.getDatanodes(), cluster.getDatanodeCount());
assertServicePortsReachable(cluster);
try (OzoneClient client = OzoneClientFactory.getRpcClient(clientConf)) {
OzoneVolume volume = client.getObjectStore().getVolume(volumeName);
OzoneBucket bucket = volume.getBucket(bucketName);
assertEquals(bucketName, bucket.getName());
+ // Key data written before the restart is still readable from the
+ // persistent datanode storage.
+ assertEquals(KEY_CONTENT, TestDataUtil.getKey(bucket, keyName));
}
}
}
diff --git a/hadoop-ozone/tools/pom.xml b/hadoop-ozone/tools/pom.xml
index bda0ebb3134..3ffab37b2b2 100644
--- a/hadoop-ozone/tools/pom.xml
+++ b/hadoop-ozone/tools/pom.xml
@@ -62,6 +62,10 @@
<groupId>org.apache.ozone</groupId>
<artifactId>hdds-config</artifactId>
</dependency>
+ <dependency>
+ <groupId>org.apache.ozone</groupId>
+ <artifactId>hdds-container-service</artifactId>
+ </dependency>
<dependency>
<groupId>org.apache.ozone</groupId>
<artifactId>hdds-server-framework</artifactId>
diff --git
a/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
b/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
index def3b5b634f..89961e8243a 100644
---
a/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
+++
b/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
@@ -17,10 +17,16 @@
package org.apache.hadoop.ozone.local;
+import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_CLIENT_ADDRESS_KEY;
+import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_CLIENT_BIND_HOST_KEY;
+import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_HTTP_ADDRESS_KEY;
+import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_HTTP_BIND_HOST_KEY;
+import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL;
import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_MIN_DATANODE;
import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_PIPELINE_CREATION;
import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_WAIT_TIME_AFTER_SAFE_MODE_EXIT;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_CONTAINER_RATIS_ENABLED_KEY;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_DATANODE_DIR_KEY;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_BLOCK_CLIENT_ADDRESS_KEY;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_BLOCK_CLIENT_BIND_HOST_KEY;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CLIENT_ADDRESS_KEY;
@@ -41,6 +47,12 @@
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_SECURITY_SERVICE_ADDRESS_KEY;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_SECURITY_SERVICE_BIND_HOST_KEY;
import static org.apache.hadoop.hdds.server.http.BaseHttpServer.SERVER_DIR;
+import static org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_IPC_PORT;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_ADMIN_PORT;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATANODE_STORAGE_DIR;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_PORT;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_IPC_PORT;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_SERVER_PORT;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_HTTP_BASEDIR;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_METADATA_DIRS;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_REPLICATION;
@@ -68,8 +80,12 @@
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.Comparator;
import java.util.HashSet;
+import java.util.List;
import java.util.Objects;
import java.util.Properties;
import java.util.Set;
@@ -85,18 +101,19 @@
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
import org.apache.hadoop.hdds.utils.IOUtils;
import org.apache.hadoop.metrics2.lib.DefaultMetricsSystem;
+import org.apache.hadoop.ozone.HddsDatanodeService;
import org.apache.hadoop.ozone.OzoneSecurityUtil;
import org.apache.hadoop.ozone.common.Storage;
+import org.apache.hadoop.ozone.container.replication.ReplicationServer;
import org.apache.hadoop.ozone.om.OMStorage;
import org.apache.hadoop.ozone.om.OzoneManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
- * Starts the SCM and OM portion of the {@code ozone local} runtime.
+ * Starts the SCM, OM, and datanode portion of the {@code ozone local} runtime.
*
- * <p>Datanodes, S3 Gateway, Recon, and end-to-end key writes are added by
- * later local runtime tickets.</p>
+ * <p>S3 Gateway and Recon are added by later local runtime tickets.</p>
*/
public final class LocalOzoneCluster implements LocalOzoneRuntime {
@@ -111,6 +128,7 @@ public final class LocalOzoneCluster implements
LocalOzoneRuntime {
private static final String OZONE_METADATA_DIR_NAME = "ozone-metadata";
private static final String DATA_DIR_NAME = "data";
private static final String RATIS_DIR_NAME = "ratis";
+ private static final String DATANODE_DIR_PREFIX = "datanode-";
private static final String SCM_CLIENT_PORT_KEY = "scm.client";
private static final String SCM_BLOCK_PORT_KEY = "scm.block";
private static final String SCM_DATANODE_PORT_KEY = "scm.datanode";
@@ -123,6 +141,15 @@ public final class LocalOzoneCluster implements
LocalOzoneRuntime {
private static final String OM_HTTP_PORT_KEY = "om.http";
private static final String OM_HTTPS_PORT_KEY = "om.https";
private static final String OM_RATIS_PORT_KEY = "om.ratis";
+ private static final String DATANODE_PORT_KEY_PREFIX = "dn.";
+ private static final String DATANODE_HTTP_PORT_KEY_SUFFIX = "http";
+ private static final String DATANODE_CLIENT_PORT_KEY_SUFFIX = "client";
+ private static final String DATANODE_CONTAINER_IPC_PORT_KEY_SUFFIX =
"container.ipc";
+ private static final String DATANODE_RATIS_IPC_PORT_KEY_SUFFIX = "ratis.ipc";
+ private static final String DATANODE_RATIS_ADMIN_PORT_KEY_SUFFIX =
"ratis.admin";
+ private static final String DATANODE_RATIS_SERVER_PORT_KEY_SUFFIX =
"ratis.server";
+ private static final String DATANODE_RATIS_DATASTREAM_PORT_KEY_SUFFIX =
"ratis.datastream";
+ private static final String DATANODE_REPLICATION_PORT_KEY_SUFFIX =
"replication";
private static final int LOCAL_RATIS_RPC_TIMEOUT_SECONDS = 1;
private static final long SCM_CLIENT_MAX_RETRY_TIMEOUT_MILLIS = 30_000;
private static final long READINESS_POLL_INTERVAL_MILLIS = 500;
@@ -144,6 +171,23 @@ public final class LocalOzoneCluster implements
LocalOzoneRuntime {
OM_RATIS_PORT_KEY
};
+ private static final String[] DATANODE_PORT_KEY_SUFFIXES = {
+ DATANODE_HTTP_PORT_KEY_SUFFIX,
+ DATANODE_CLIENT_PORT_KEY_SUFFIX,
+ DATANODE_CONTAINER_IPC_PORT_KEY_SUFFIX,
+ DATANODE_RATIS_IPC_PORT_KEY_SUFFIX,
+ DATANODE_RATIS_ADMIN_PORT_KEY_SUFFIX,
+ DATANODE_RATIS_SERVER_PORT_KEY_SUFFIX,
+ DATANODE_RATIS_DATASTREAM_PORT_KEY_SUFFIX,
+ DATANODE_REPLICATION_PORT_KEY_SUFFIX
+ };
+
+ // Every datanode runs in this JVM and reserves
DATANODE_PORT_KEY_SUFFIXES.length
+ // local ports, so an unbounded count would exhaust local ports; cap it.
+ static final int MAX_DATANODES = 20;
+
+ private static final String[] NO_ARGS = new String[0];
+
private final LocalOzoneClusterConfig config;
private final OzoneConfiguration seedConfiguration;
private boolean closed;
@@ -151,6 +195,7 @@ public final class LocalOzoneCluster implements
LocalOzoneRuntime {
private PreparedConfiguration preparedConfiguration;
private StorageContainerManager scm;
private OzoneManager om;
+ private final List<HddsDatanodeService> datanodes = new ArrayList<>();
private boolean previousMetricsMiniClusterMode;
private boolean metricsMiniClusterModeEnabled;
@@ -179,7 +224,8 @@ public void start() throws Exception {
initializeStorage(prepared.getConfiguration());
startScm(prepared.getConfiguration());
startOm(prepared.getConfiguration());
- waitForScmAndOmReadiness(config.getStartupTimeout());
+ startDatanodes(prepared.getDatanodeConfigurations());
+ waitForClusterReadiness(config.getStartupTimeout());
} catch (Exception ex) {
// Roll back without latching closed: the caller's close() still owns the
// ephemeral data dir lifecycle.
@@ -202,9 +248,12 @@ PreparedConfiguration prepareConfiguration() throws
IOException {
PortAllocator portAllocator = new PortAllocator();
int scmPort = configureScm(conf, persistedPorts, portAllocator);
int omPort = configureOm(conf, persistedPorts, portAllocator);
+ List<OzoneConfiguration> datanodeConfigurations =
+ configureDatanodes(conf, persistedPorts, portAllocator);
persistedPorts.store();
- preparedConfiguration = new PreparedConfiguration(conf, scmPort, omPort);
+ preparedConfiguration = new PreparedConfiguration(conf, scmPort, omPort,
+ datanodeConfigurations);
return preparedConfiguration;
}
@@ -224,6 +273,13 @@ public int getOmPort() {
return om.getOmRpcServerAddr().getPort();
}
+ /**
+ * Returns the number of running datanodes.
+ */
+ public int getDatanodeCount() {
+ return datanodes.size();
+ }
+
@Override
public int getS3gPort() {
return -1;
@@ -258,13 +314,26 @@ public void close() throws IOException {
private void stopServices() {
try {
- // Shutdown is best-effort so one failed service cannot leak the other.
- IOUtils.closeQuietly(this::stopOm, this::stopScm);
+ // Shutdown is best-effort so one failed service cannot leak the others.
+ IOUtils.closeQuietly(this::stopDatanodes, this::stopOm, this::stopScm);
} finally {
restoreSameJvmMetricsMode();
}
}
+ private void stopDatanodes() {
+ List<AutoCloseable> stoppers = new ArrayList<>();
+ for (int i = datanodes.size() - 1; i >= 0; i--) {
+ HddsDatanodeService service = datanodes.get(i);
+ stoppers.add(() -> {
+ service.stop();
+ service.join();
+ });
+ }
+ datanodes.clear();
+ IOUtils.closeQuietly(stoppers);
+ }
+
private void stopOm() {
OzoneManager service = om;
om = null;
@@ -293,6 +362,11 @@ private void configureLocalDefaults(OzoneConfiguration
conf) {
conf.set(OZONE_SERVER_DEFAULT_REPLICATION_TYPE_KEY,
ReplicationType.STAND_ALONE.name());
conf.setBoolean(HDDS_CONTAINER_RATIS_ENABLED_KEY, false);
+ // A single-node local cluster can heartbeat aggressively; this speeds
+ // datanode registration and safe-mode exit. Use set(), not setIfUnset():
+ // ozone-default.xml supplies the 30s default that would otherwise defeat
+ // the override.
+ conf.set(HDDS_HEARTBEAT_INTERVAL, "1s");
conf.setBoolean(HDDS_SCM_SAFEMODE_PIPELINE_CREATION, false);
conf.setInt(HDDS_SCM_SAFEMODE_MIN_DATANODE,
Math.max(1, config.getDatanodes()));
@@ -403,6 +477,81 @@ private void configureOmStorage(OzoneConfiguration conf)
throws IOException {
omMetadataDir.toString());
}
+ private List<OzoneConfiguration> configureDatanodes(OzoneConfiguration conf,
+ PersistedPortState persistedPorts, PortAllocator portAllocator)
+ throws IOException {
+ int datanodeCount = config.getDatanodes();
+ if (datanodeCount > MAX_DATANODES) {
+ throw new IOException("Datanode count " + datanodeCount
+ + " exceeds the local maximum of " + MAX_DATANODES
+ + "; each datanode reserves " + DATANODE_PORT_KEY_SUFFIXES.length
+ + " local ports.");
+ }
+ List<OzoneConfiguration> datanodeConfigurations =
+ new ArrayList<>(config.getDatanodes());
+ for (int index = 0; index < config.getDatanodes(); index++) {
+ datanodeConfigurations.add(
+ configureDatanode(conf, index, persistedPorts, portAllocator));
+ }
+ return datanodeConfigurations;
+ }
+
+ private OzoneConfiguration configureDatanode(OzoneConfiguration conf,
+ int index, PersistedPortState persistedPorts,
+ PortAllocator portAllocator) throws IOException {
+ OzoneConfiguration dnConf = new OzoneConfiguration(conf);
+ configureDatanodeStorage(dnConf, index);
+
+ dnConf.set(HDDS_DATANODE_HTTP_ADDRESS_KEY, address(config.getHost(),
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_HTTP_PORT_KEY_SUFFIX)));
+ dnConf.set(HDDS_DATANODE_HTTP_BIND_HOST_KEY, config.getBindHost());
+ dnConf.set(HDDS_DATANODE_CLIENT_ADDRESS_KEY, address(config.getHost(),
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_CLIENT_PORT_KEY_SUFFIX)));
+ dnConf.set(HDDS_DATANODE_CLIENT_BIND_HOST_KEY, config.getBindHost());
+ dnConf.setInt(HDDS_CONTAINER_IPC_PORT,
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_CONTAINER_IPC_PORT_KEY_SUFFIX));
+ dnConf.setInt(HDDS_CONTAINER_RATIS_IPC_PORT,
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_RATIS_IPC_PORT_KEY_SUFFIX));
+ dnConf.setInt(HDDS_CONTAINER_RATIS_ADMIN_PORT,
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_RATIS_ADMIN_PORT_KEY_SUFFIX));
+ dnConf.setInt(HDDS_CONTAINER_RATIS_SERVER_PORT,
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_RATIS_SERVER_PORT_KEY_SUFFIX));
+ dnConf.setInt(HDDS_CONTAINER_RATIS_DATASTREAM_PORT,
+ reserveDatanodePort(portAllocator, persistedPorts, index,
+ DATANODE_RATIS_DATASTREAM_PORT_KEY_SUFFIX));
+
+ ReplicationServer.ReplicationConfig replicationConfig =
+ dnConf.getObject(ReplicationServer.ReplicationConfig.class);
+ replicationConfig.setPort(reserveDatanodePort(portAllocator,
+ persistedPorts, index, DATANODE_REPLICATION_PORT_KEY_SUFFIX));
+ dnConf.setFromObject(replicationConfig);
+ return dnConf;
+ }
+
+ private void configureDatanodeStorage(OzoneConfiguration dnConf, int index)
+ throws IOException {
+ Path datanodeDir = config.getDataDir()
+ .resolve(DATANODE_DIR_PREFIX + index);
+ Path datanodeMetadataDir = datanodeDir.resolve(OZONE_METADATA_DIR_NAME);
+ Files.createDirectories(datanodeMetadataDir);
+ Files.createDirectories(datanodeDir.resolve(DATA_DIR_NAME));
+
+ dnConf.set(OZONE_METADATA_DIRS, datanodeMetadataDir.toString());
+ // Each datanode gets its own Jetty base dir so same-JVM HTTP servers do
+ // not share unpacked web resources.
+ dnConf.set(OZONE_HTTP_BASEDIR, datanodeMetadataDir + SERVER_DIR);
+ dnConf.set(HDDS_DATANODE_DIR_KEY,
+ datanodeDir.resolve(DATA_DIR_NAME).toString());
+ dnConf.set(HDDS_CONTAINER_RATIS_DATANODE_STORAGE_DIR,
+ datanodeDir.resolve(RATIS_DIR_NAME).toString());
+ }
+
private void initializeStorage(OzoneConfiguration conf) throws IOException {
SCMStorageConfig scmStorage = new SCMStorageConfig(conf);
OMStorage omStorage = new OMStorage(conf);
@@ -483,22 +632,43 @@ private void startOm(OzoneConfiguration conf) throws
Exception {
om.start();
}
- private void waitForScmAndOmReadiness(Duration timeout) throws Exception {
+ private void startDatanodes(List<OzoneConfiguration> datanodeConfigurations)
{
+ for (OzoneConfiguration dnConf : datanodeConfigurations) {
+ // Track the datanode before start() so a failed start can still be
+ // rolled back by stopServices().
+ HddsDatanodeService datanode = new HddsDatanodeService(NO_ARGS);
+ datanodes.add(datanode);
+ datanode.start(dnConf);
+ }
+ }
+
+ private void waitForClusterReadiness(Duration timeout) throws Exception {
long deadlineNanos = System.nanoTime() + timeout.toNanos();
while (true) {
- // Readiness is intentionally scoped to SCM and OM leadership until
- // datanode and full safe-mode readiness are added in later tickets.
- if (scm.checkLeader() && om.isLeaderReady()) {
+ if (isClusterReady()) {
return;
}
if (System.nanoTime() >= deadlineNanos) {
throw new TimeoutException("Timed out waiting " + timeout
- + " for local SCM and OM leadership.");
+ + " for the local Ozone cluster to become ready.");
}
Thread.sleep(READINESS_POLL_INTERVAL_MILLIS);
}
}
+ private boolean isClusterReady() {
+ if (!scm.checkLeader() || !om.isLeaderReady()) {
+ return false;
+ }
+ if (config.getDatanodes() == 0) {
+ return true;
+ }
+ // The cluster is usable once every datanode has registered with SCM and
+ // SCM has left safe mode.
+ return scm.getScmNodeManager().getAllNodes().size() >=
config.getDatanodes()
+ && !scm.isInSafeMode();
+ }
+
private void enableSameJvmMetricsMode() {
if (!metricsMiniClusterModeEnabled) {
previousMetricsMiniClusterMode =
DefaultMetricsSystem.inMiniClusterMode();
@@ -567,11 +737,22 @@ private PersistedPortState loadPersistedPortState()
throws IOException {
PersistedPortState persistedPorts =
PersistedPortState.load(portStateFile());
if (config.getFormatMode() == LocalOzoneClusterConfig.FormatMode.NEVER) {
- persistedPorts.requireKeys(REQUIRED_PERSISTED_PORT_KEYS);
+ persistedPorts.requireKeys(requiredPersistedPortKeys());
}
return persistedPorts;
}
+ private String[] requiredPersistedPortKeys() {
+ List<String> keys =
+ new ArrayList<>(Arrays.asList(REQUIRED_PERSISTED_PORT_KEYS));
+ for (int index = 0; index < config.getDatanodes(); index++) {
+ for (String suffix : DATANODE_PORT_KEY_SUFFIXES) {
+ keys.add(datanodePortKey(index, suffix));
+ }
+ }
+ return keys.toArray(new String[0]);
+ }
+
private int reservePort(PortAllocator allocator,
PersistedPortState persistedPorts, String key, int configuredPort)
throws IOException {
@@ -582,6 +763,17 @@ private int reservePort(PortAllocator allocator,
return port;
}
+ private int reserveDatanodePort(PortAllocator allocator,
+ PersistedPortState persistedPorts, int index, String suffix)
+ throws IOException {
+ return reservePort(allocator, persistedPorts,
+ datanodePortKey(index, suffix), 0);
+ }
+
+ private static String datanodePortKey(int index, String suffix) {
+ return DATANODE_PORT_KEY_PREFIX + index + "." + suffix;
+ }
+
private Path metadataDir() {
return config.getDataDir().resolve(METADATA_DIR_NAME);
}
@@ -615,13 +807,17 @@ static final class PreparedConfiguration {
private final OzoneConfiguration configuration;
private final int scmPort;
private final int omPort;
+ private final List<OzoneConfiguration> datanodeConfigurations;
PreparedConfiguration(OzoneConfiguration configuration, int scmPort,
- int omPort) {
+ int omPort, List<OzoneConfiguration> datanodeConfigurations) {
this.configuration = Objects.requireNonNull(configuration,
"configuration");
this.scmPort = scmPort;
this.omPort = omPort;
+ this.datanodeConfigurations = Collections.unmodifiableList(
+ new ArrayList<>(Objects.requireNonNull(datanodeConfigurations,
+ "datanodeConfigurations")));
}
OzoneConfiguration getConfiguration() {
@@ -635,6 +831,10 @@ int getScmPort() {
int getOmPort() {
return omPort;
}
+
+ List<OzoneConfiguration> getDatanodeConfigurations() {
+ return datanodeConfigurations;
+ }
}
/**
diff --git
a/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
b/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
index 8ecc3a9db98..1f109c6de42 100644
---
a/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
+++
b/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
@@ -21,8 +21,10 @@
import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_MIN_DATANODE;
import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_PIPELINE_CREATION;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_CONTAINER_RATIS_ENABLED_KEY;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_DATANODE_DIR_KEY;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CLIENT_ADDRESS_KEY;
import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_NAMES;
+import static org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_IPC_PORT;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_METADATA_DIRS;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_REPLICATION;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_REPLICATION_TYPE;
@@ -305,6 +307,113 @@ void persistedPortFileContainsDistinctAllocatedPorts()
throws Exception {
properties.getProperty("om.rpc"));
}
+ @Test
+ void prepareConfigurationCreatesDatanodeConfigurations() throws Exception {
+ Path dataDir = tempDir.resolve("local-ozone");
+ LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
+ .setDatanodes(2)
+ .build();
+
+ LocalOzoneCluster.PreparedConfiguration prepared = prepare(config);
+
+ assertEquals(2, prepared.getDatanodeConfigurations().size());
+ for (int index = 0; index < 2; index++) {
+ OzoneConfiguration dnConf =
prepared.getDatanodeConfigurations().get(index);
+ Path datanodeDir = dataDir.resolve("datanode-" + index);
+ assertTrue(Files.isDirectory(datanodeDir.resolve("ozone-metadata")));
+ assertTrue(Files.isDirectory(datanodeDir.resolve("data")));
+ assertEquals(datanodeDir.resolve("ozone-metadata").toString(),
+ dnConf.get(OZONE_METADATA_DIRS));
+ assertEquals(datanodeDir.resolve("data").toString(),
+ dnConf.get(HDDS_DATANODE_DIR_KEY));
+ assertTrue(dnConf.getInt(HDDS_CONTAINER_IPC_PORT, 0) > 0);
+ }
+ assertNotEquals(
+ prepared.getDatanodeConfigurations().get(0)
+ .getInt(HDDS_CONTAINER_IPC_PORT, 0),
+ prepared.getDatanodeConfigurations().get(1)
+ .getInt(HDDS_CONTAINER_IPC_PORT, 0));
+ }
+
+ @Test
+ void prepareConfigurationRejectsTooManyDatanodes() throws Exception {
+ LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(
+ tempDir.resolve("local-ozone"))
+ .setDatanodes(LocalOzoneCluster.MAX_DATANODES + 1)
+ .build();
+
+ IOException error = assertPrepareFails(config);
+
+ assertEquals("Datanode count " + (LocalOzoneCluster.MAX_DATANODES + 1)
+ + " exceeds the local maximum of " + LocalOzoneCluster.MAX_DATANODES
+ + "; each datanode reserves 8 local ports.", error.getMessage());
+ }
+
+ @Test
+ void persistedPortFileContainsDatanodePorts() throws Exception {
+ Path dataDir = tempDir.resolve("local-ozone");
+ LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
+ .setDatanodes(1)
+ .build();
+
+ prepare(config);
+
+ Properties properties = loadPortState(dataDir);
+ assertPositivePort(properties, "dn.0.http");
+ assertPositivePort(properties, "dn.0.client");
+ assertPositivePort(properties, "dn.0.container.ipc");
+ assertPositivePort(properties, "dn.0.ratis.ipc");
+ assertPositivePort(properties, "dn.0.ratis.admin");
+ assertPositivePort(properties, "dn.0.ratis.server");
+ assertPositivePort(properties, "dn.0.ratis.datastream");
+ assertPositivePort(properties, "dn.0.replication");
+ }
+
+ @Test
+ void prepareConfigurationPersistsDatanodePortsAcrossInstances()
+ throws Exception {
+ Path dataDir = tempDir.resolve("local-ozone");
+ LocalOzoneClusterConfig config =
+ LocalOzoneClusterConfig.builder(dataDir).build();
+
+ LocalOzoneCluster.PreparedConfiguration first = prepare(config);
+ LocalOzoneCluster.PreparedConfiguration second = prepare(config);
+
+ assertEquals(
+ first.getDatanodeConfigurations().get(0)
+ .getInt(HDDS_CONTAINER_IPC_PORT, 0),
+ second.getDatanodeConfigurations().get(0)
+ .getInt(HDDS_CONTAINER_IPC_PORT, 0));
+ }
+
+ @Test
+ void formatNeverRejectsPortStateMissingDatanodePorts() throws Exception {
+ Path dataDir = tempDir.resolve("local-ozone");
+ prepare(LocalOzoneClusterConfig.builder(dataDir).setDatanodes(1).build());
+ LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
+ .setFormatMode(LocalOzoneClusterConfig.FormatMode.NEVER)
+ .setDatanodes(2)
+ .build();
+
+ IOException error = assertPrepareFails(config);
+
+ assertMessageContains(error, "dn.1.");
+ }
+
+ @Test
+ void getDatanodeCountReturnsZeroBeforeStart() throws Exception {
+ LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(
+ tempDir.resolve("local-ozone"))
+ .setDatanodes(3)
+ .build();
+
+ try (LocalOzoneCluster cluster = newCluster(config)) {
+ cluster.prepareConfiguration();
+
+ assertEquals(0, cluster.getDatanodeCount());
+ }
+ }
+
private LocalOzoneCluster.PreparedConfiguration prepare(
LocalOzoneClusterConfig config) throws IOException {
try (LocalOzoneCluster cluster = newCluster(config)) {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]