This is an automated email from the ASF dual-hosted git repository.
szetszwo 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 5197855a329 HDDS-12991. Recreate pipelines after enabling ratis write
streaming (#10842)
5197855a329 is described below
commit 5197855a329cef5db4509c62f4beffc14bbf4756
Author: Andrey Yarovoy <[email protected]>
AuthorDate: Fri Jul 31 14:55:04 2026 -0400
HDDS-12991. Recreate pipelines after enabling ratis write streaming (#10842)
---
.../hadoop/hdds/protocol/DatanodeDetails.java | 21 +-
.../hadoop/hdds/scm/node/SCMNodeManager.java | 5 +
.../hdds/scm/pipeline/PipelineManagerImpl.java | 68 +++-
.../hadoop/hdds/scm/node/TestSCMNodeManager.java | 77 ++++
.../hdds/scm/pipeline/TestPipelineManagerImpl.java | 186 +++++++++
.../TestOzoneFileSystemDataStreamEnablement.java | 427 +++++++++++++++++++++
6 files changed, 776 insertions(+), 8 deletions(-)
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
index abd00e70683..4abe44a3040 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
@@ -35,6 +35,7 @@
import java.util.Objects;
import java.util.Set;
import java.util.UUID;
+import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
import org.apache.hadoop.hdds.DatanodeVersion;
import org.apache.hadoop.hdds.HddsUtils;
@@ -384,9 +385,27 @@ public synchronized boolean hasPort(Port.Name name) {
return false;
}
+ /**
+ * Whether this datanode's exposed ports differ from {@code other}'s.
+ * Compared as a set of name=value entries, since {@link Port#equals}
+ * ignores the port value.
+ *
+ * @param other another snapshot of this datanode
+ * @return true if the two port sets are not identical
+ */
+ public boolean portsChanged(DatanodeDetails other) {
+ return !portValues(this).equals(portValues(other));
+ }
+
+ private static Set<String> portValues(DatanodeDetails datanodeDetails) {
+ return datanodeDetails.getPorts().stream()
+ .map(port -> port.getName() + "=" + port.getValue())
+ .collect(Collectors.toSet());
+ }
+
/**
* Helper method to get the Ratis port.
- *
+ *
* @return Port
*/
public Port getRatisPort() {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
index cbcfc553791..fe886bf0c20 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
@@ -474,6 +474,11 @@ public RegisteredCommand register(
"oldVersion = {}, newVersion = {}.",
datanodeDetails, oldNode.getVersion(),
datanodeDetails.getVersion());
nodeStateManager.updateNode(datanodeDetails, layoutInfo);
+ } else if (oldNode.portsChanged(datanodeDetails)) {
+ // Refresh the stored node when its port set changes
+ LOG.info("Updating ports for registered datanode {}: {} -> {}",
+ datanodeDetails, oldNode.getPorts(), datanodeDetails.getPorts());
+ nodeStateManager.updateNode(datanodeDetails, layoutInfo);
}
} catch (NodeNotFoundException e) {
LOG.error("Cannot find datanode {} from nodeStateManager",
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
index 8003d8d52e9..c472e65b8b3 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
@@ -63,6 +63,7 @@
import org.apache.hadoop.hdds.utils.db.Table;
import org.apache.hadoop.metrics2.util.MBeans;
import org.apache.hadoop.ozone.ClientVersion;
+import org.apache.hadoop.ozone.OzoneConfigKeys;
import org.apache.hadoop.util.Time;
import org.apache.ratis.protocol.exceptions.NotLeaderException;
import org.slf4j.Logger;
@@ -182,13 +183,8 @@ public static PipelineManagerImpl newPipelineManager(
.setServiceName("BackgroundPipelineScrubber")
.setIntervalInMillis(scrubberIntervalInMillis)
.setWaitTimeInMillis(safeModeWaitMs)
- .setPeriodicalTask(() -> {
- try {
- pipelineManager.scrubPipelines();
- } catch (IOException e) {
- LOG.error("Unexpected error during pipeline scrubbing", e);
- }
- }).build();
+
.setPeriodicalTask(pipelineManager::scrubAndClosePipelinesMissingDataStreamPort)
+ .build();
pipelineManager.setBackgroundPipelineScrubber(backgroundPipelineScrubber);
serviceManager.register(backgroundPipelineScrubber);
@@ -562,6 +558,64 @@ static boolean
sameIdDifferentHostOrAddress(DatanodeDetails left, DatanodeDetail
|| !left.getHostName().equals(right.getHostName()));
}
+ /**
+ * Scrub pipelines, then close (and delete) OPEN RATIS pipelines whose
+ * registered nodes now advertise the RATIS_DATASTREAM port their stored node
+ * snapshot lacks.
+ */
+ public void scrubAndClosePipelinesMissingDataStreamPort() {
+ try {
+ scrubPipelines();
+ } catch (IOException e) {
+ LOG.error("Unexpected error during pipeline scrubbing", e);
+ }
+ closePipelinesMissingDataStreamPort();
+ }
+
+ void closePipelinesMissingDataStreamPort() {
+ if (!isDataStreamEnabled()) {
+ return;
+ }
+ for (Pipeline pipeline : getPipelines()) {
+ if (!pipeline.isOpen()
+ || pipeline.getType() != ReplicationType.RATIS
+ || !nodesMissingDataStreamPort(pipeline)) {
+ continue;
+ }
+ try {
+ final PipelineID id = pipeline.getId();
+ LOG.info("Closing RATIS pipeline {}", id);
+ closePipeline(id);
+ deletePipeline(id);
+ } catch (IOException e) {
+ LOG.error("Failed to close RATIS pipeline {} missing the datastream "
+ + "port", pipeline.getId(), e);
+ }
+ }
+ }
+
+ private boolean isDataStreamEnabled() {
+ return conf.getBoolean(
+ OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED,
+ OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED_DEFAULT);
+ }
+
+ /**
+ * Whether any registered node of the pipeline now advertises the
+ * RATIS_DATASTREAM port that the pipeline's stored node snapshot lacks.
+ */
+ private boolean nodesMissingDataStreamPort(Pipeline pipeline) {
+ for (DatanodeDetails stored : pipeline.getNodes()) {
+ final DatanodeDetails current = nodeManager.getNode(stored.getID());
+ if (current != null
+ && current.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM)
+ && !stored.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
/**
* Scrub pipelines.
*/
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
index a0f1c1bf362..78814b53e81 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
@@ -43,6 +43,7 @@
import static org.apache.ozone.test.MetricsAsserts.getMetrics;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -63,6 +64,7 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.UUID;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
@@ -344,6 +346,81 @@ private DatanodeDetails
registerWithCapacity(SCMNodeManager nodeManager,
return cmd.getDatanode();
}
+ private static DatanodeDetails.Builder datanodeWithoutDatastream(UUID uuid) {
+ return DatanodeDetails.newBuilder()
+ .setUuid(uuid)
+ .setHostName("host-" + uuid)
+ .setIpAddress("127.0.0.1")
+ .addPort(DatanodeDetails.newPort(
+ DatanodeDetails.Port.Name.STANDALONE, 9859))
+ .addPort(DatanodeDetails.newPort(
+ DatanodeDetails.Port.Name.RATIS, 9858));
+ }
+
+ private void registerNode(SCMNodeManager nodeManager, DatanodeDetails dn) {
+ StorageReportProto storageReport = HddsTestUtils.createStorageReport(
+ dn.getID(), dn.getNetworkFullPath(), Long.MAX_VALUE);
+ MetadataStorageReportProto metadataStorageReport =
+ HddsTestUtils.createMetadataStorageReport(
+ dn.getNetworkFullPath(), Long.MAX_VALUE);
+ RegisteredCommand cmd = nodeManager.register(dn,
+ HddsTestUtils.createNodeReport(Arrays.asList(storageReport),
+ Arrays.asList(metadataStorageReport)),
+ getRandomPipelineReports(), UpgradeUtils.defaultLayoutVersionProto());
+ assertEquals(success, cmd.getError());
+ }
+
+ /**
+ * A datanode that re-registers with the same identity but now exposes the
+ * RATIS_DATASTREAM port (e.g. Ratis DataStream was enabled) must have its
+ * stored record refreshed so the new port is visible (HDDS-15799).
+ */
+ @Test
+ public void testRegisterRefreshesPortsOnPortChange()
+ throws IOException, AuthenticationException {
+ try (SCMNodeManager nodeManager = createNodeManager(getConf())) {
+ final UUID uuid = UUID.randomUUID();
+
+ // First registration: streaming disabled, no RATIS_DATASTREAM port.
+ registerNode(nodeManager, datanodeWithoutDatastream(uuid).build());
+ DatanodeDetails stored = nodeManager.getNode(
+ datanodeWithoutDatastream(uuid).build().getID());
+ assertFalse(stored.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM));
+
+ // Re-registration (same id/ip/host/version) now exposing the port.
+ registerNode(nodeManager, datanodeWithoutDatastream(uuid)
+ .addPort(DatanodeDetails.newPort(
+ DatanodeDetails.Port.Name.RATIS_DATASTREAM, 9855))
+ .build());
+ stored = nodeManager.getNode(
+ datanodeWithoutDatastream(uuid).build().getID());
+ assertTrue(stored.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM),
+ "stored node should be refreshed with the RATIS_DATASTREAM port");
+ }
+ }
+
+ /**
+ * Re-registering a datanode with an unchanged port set must not disturb the
+ * stored record (the port-refresh branch is skipped).
+ */
+ @Test
+ public void testRegisterKeepsPortsWhenUnchanged()
+ throws IOException, AuthenticationException {
+ try (SCMNodeManager nodeManager = createNodeManager(getConf())) {
+ final UUID uuid = UUID.randomUUID();
+
+ registerNode(nodeManager, datanodeWithoutDatastream(uuid).build());
+ // Re-register with the identical port set.
+ registerNode(nodeManager, datanodeWithoutDatastream(uuid).build());
+
+ final DatanodeDetails stored = nodeManager.getNode(
+ datanodeWithoutDatastream(uuid).build().getID());
+ assertTrue(stored.hasPort(DatanodeDetails.Port.Name.STANDALONE));
+ assertTrue(stored.hasPort(DatanodeDetails.Port.Name.RATIS));
+ assertFalse(stored.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM));
+ }
+ }
+
private void assertPipelineClosedAfterLayoutHeartbeat(
DatanodeDetails originalNode1, DatanodeDetails originalNode2,
SCMNodeManager nodeManager, LayoutVersionProto layout) throws Exception {
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java
index 2f1d5ede5d7..e980b81318d 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java
@@ -23,6 +23,7 @@
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_PIPELINE_DESTROY_TIMEOUT;
import static
org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState.ALLOCATED;
import static org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState.OPEN;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED;
import static org.apache.ozone.test.MetricsAsserts.getLongCounter;
import static org.apache.ozone.test.MetricsAsserts.getMetrics;
import static org.apache.ratis.util.Preconditions.assertInstanceOf;
@@ -88,6 +89,7 @@
import org.apache.hadoop.hdds.scm.ha.SCMHAManagerStub;
import org.apache.hadoop.hdds.scm.ha.SCMServiceManager;
import org.apache.hadoop.hdds.scm.metadata.SCMDBDefinition;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.node.NodeStatus;
import
org.apache.hadoop.hdds.scm.pipeline.choose.algorithms.HealthyPipelineChoosePolicy;
@@ -100,6 +102,7 @@
import org.apache.hadoop.hdds.utils.db.DBStoreBuilder;
import org.apache.hadoop.hdds.utils.db.Table;
import org.apache.hadoop.metrics2.MetricsRecordBuilder;
+import org.apache.hadoop.ozone.ClientVersion;
import org.apache.hadoop.ozone.container.common.SCMTestUtils;
import org.apache.ozone.test.GenericTestUtils;
import org.apache.ozone.test.GenericTestUtils.LogCapturer;
@@ -980,4 +983,187 @@ private static void
assertFailsNotLeader(CheckedRunnable<?> block) {
assertEquals(ResultCodes.SCM_NOT_LEADER, e.getResult());
assertInstanceOf(NotLeaderException.class, e.getCause());
}
+
+ private static DatanodeDetails portlessDatanode(DatanodeID id) {
+ return DatanodeDetails.newBuilder()
+ .setID(id)
+ .setHostName("host-" + id)
+ .setIpAddress("127.0.0.1")
+ .addPort(DatanodeDetails.newPort(
+ DatanodeDetails.Port.Name.STANDALONE, 9859))
+ .addPort(DatanodeDetails.newPort(
+ DatanodeDetails.Port.Name.RATIS, 9858))
+ .build();
+ }
+
+ private Pipeline addPipeline(PipelineManagerImpl pipelineManager,
+ Pipeline.PipelineState state, List<DatanodeDetails> nodes)
+ throws IOException {
+ final Pipeline pipeline = Pipeline.newBuilder()
+ .setReplicationConfig(
+ RatisReplicationConfig.getInstance(ReplicationFactor.THREE))
+ .setNodes(nodes)
+ .setState(state)
+ .setId(PipelineID.randomId())
+ .build();
+ pipelineManager.getStateManager().addPipeline(
+ pipeline.getProtobufMessage(ClientVersion.CURRENT_VERSION));
+ return pipeline;
+ }
+
+ private static boolean exists(PipelineManagerImpl pipelineManager,
+ PipelineID id) {
+ try {
+ pipelineManager.getPipeline(id);
+ return true;
+ } catch (PipelineNotFoundException e) {
+ return false;
+ }
+ }
+
+ @Test
+ public void testClosePipelinesExposingNewPorts() throws Exception {
+ conf.setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ // Registered datanodes (MockNodeManager) expose all ports incl
datastream.
+ final List<DatanodeInfo> registered = nodeManager.getAllNodes();
+ final List<DatanodeDetails> idsA = new ArrayList<>();
+ final List<DatanodeDetails> idsB = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ idsA.add(portlessDatanode(registered.get(i).getID()));
+ idsB.add(portlessDatanode(registered.get(i + 3).getID()));
+ }
+
+ // OPEN, registered nodes, portless -> legacy pipeline, must be closed.
+ final Pipeline stale = addPipeline(pipelineManager, OPEN, idsA);
+ // OPEN, registered nodes carrying all ports -> not stale, kept.
+ final Pipeline portful = addPipeline(pipelineManager, OPEN,
+ new ArrayList<>(registered.subList(6, 9)));
+ // OPEN, but nodes are NOT registered -> cannot heal, left alone.
+ final List<DatanodeDetails> unregistered = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ unregistered.add(portlessDatanode(DatanodeID.randomID()));
+ }
+ final Pipeline unreg = addPipeline(pipelineManager, OPEN, unregistered);
+ // ALLOCATED (non-open) portless -> skipped.
+ final Pipeline allocated = addPipeline(pipelineManager, ALLOCATED, idsB);
+
+ pipelineManager.closePipelinesMissingDataStreamPort();
+
+ assertFalse(exists(pipelineManager, stale.getId()),
+ "OPEN pipeline whose nodes expose new ports should be closed and
deleted");
+ assertTrue(exists(pipelineManager, portful.getId()));
+ assertTrue(exists(pipelineManager, unreg.getId()));
+ assertTrue(exists(pipelineManager, allocated.getId()));
+ }
+ }
+
+ @Test
+ public void testClosePipelinesExposingNewPortsSkippedWhenDataStreamDisabled()
+ throws Exception {
+ // Datastream disabled (default): even a portless RATIS pipeline is kept.
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ final List<DatanodeInfo> registered = nodeManager.getAllNodes();
+ final List<DatanodeDetails> nodes = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ nodes.add(portlessDatanode(registered.get(i).getID()));
+ }
+ final Pipeline portless = addPipeline(pipelineManager, OPEN, nodes);
+
+ pipelineManager.closePipelinesMissingDataStreamPort();
+
+ assertTrue(exists(pipelineManager, portless.getId()),
+ "portless pipeline must be kept while datastream is disabled");
+ }
+ }
+
+ @Test
+ public void testClosePipelinesExposingNewPortsSkipsEcPipeline()
+ throws Exception {
+ conf.setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ final List<DatanodeInfo> registered = nodeManager.getAllNodes();
+ final List<DatanodeDetails> nodes = new ArrayList<>();
+ for (int i = 0; i < 5; i++) {
+ nodes.add(portlessDatanode(registered.get(i).getID()));
+ }
+ final Pipeline ec = Pipeline.newBuilder()
+ .setReplicationConfig(new ECReplicationConfig(3, 2))
+ .setNodes(nodes)
+ .setState(OPEN)
+ .setId(PipelineID.randomId())
+ .build();
+ pipelineManager.getStateManager().addPipeline(
+ ec.getProtobufMessage(ClientVersion.CURRENT_VERSION));
+
+ pipelineManager.closePipelinesMissingDataStreamPort();
+
+ assertTrue(exists(pipelineManager, ec.getId()),
+ "EC pipeline must not be closed by datastream port scrubbing");
+ }
+ }
+
+ @Test
+ public void testClosePipelinesExposingNewPortsKeepsNotYetRestartedNodes()
+ throws Exception {
+ conf.setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ // Nodes are registered and healthy but still lack the datastream port
+ // (they have not restarted yet during a rolling enablement).
+ final List<DatanodeDetails> nodes = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ final DatanodeDetails portless =
portlessDatanode(DatanodeID.randomID());
+ nodeManager.register(new DatanodeInfo(portless,
+ NodeStatus.inServiceHealthy(), null,
+ HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT), null, null);
+ nodes.add(portless);
+ }
+ final Pipeline pending = addPipeline(pipelineManager, OPEN, nodes);
+
+ pipelineManager.closePipelinesMissingDataStreamPort();
+
+ assertTrue(exists(pipelineManager, pending.getId()),
+ "pipeline whose registered nodes have not yet advertised the "
+ + "datastream port must be kept");
+ }
+ }
+
+ @Test
+ public void testClosePipelinesExposingNewPortsSwallowsError() throws
Exception {
+ conf.setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ final List<DatanodeInfo> registered = nodeManager.getAllNodes();
+ final List<DatanodeDetails> nodes = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ nodes.add(portlessDatanode(registered.get(i).getID()));
+ }
+ final Pipeline stale = addPipeline(pipelineManager, OPEN, nodes);
+
+ final PipelineManagerImpl spy = spy(pipelineManager);
+ doThrow(new IOException("boom")).when(spy).closePipeline(stale.getId());
+ // The close failure is logged and swallowed; the loop does not throw.
+ spy.closePipelinesMissingDataStreamPort();
+ assertTrue(exists(pipelineManager, stale.getId()));
+ }
+ }
+
+ @Test
+ public void testScrubAndCloseWiring() throws Exception {
+ // The background task scrubs then closes pipelines exposing new ports; on
+ // an empty manager both are no-ops and must not throw.
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ pipelineManager.scrubAndClosePipelinesMissingDataStreamPort();
+ }
+ }
+
+ @Test
+ public void testScrubAndCloseSwallowsScrubError() throws Exception {
+ try (PipelineManagerImpl pipelineManager = createPipelineManager(true)) {
+ final PipelineManagerImpl spy = spy(pipelineManager);
+ doThrow(new IOException("boom")).when(spy).scrubPipelines();
+ // Scrub failure is logged and swallowed; the close pass still runs.
+ spy.scrubAndClosePipelinesMissingDataStreamPort();
+ verify(spy).closePipelinesMissingDataStreamPort();
+ }
+ }
}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java
new file mode 100644
index 00000000000..11a55b471bd
--- /dev/null
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java
@@ -0,0 +1,427 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hadoop.fs.ozone;
+
+import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL;
+import static
org.apache.hadoop.hdds.protocol.DatanodeDetails.Port.Name.RATIS_DATASTREAM;
+import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
+import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL;
+import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL;
+import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_PIPELINE_CREATION_INTERVAL;
+import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_PIPELINE_SCRUB_INTERVAL;
+import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_RATIS_PIPELINE_LIMIT;
+import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_FS_DATASTREAM_AUTO_THRESHOLD;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_FS_DATASTREAM_ENABLED;
+import static org.apache.hadoop.ozone.OzoneConsts.OZONE_URI_SCHEME;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.IOException;
+import java.util.List;
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.function.BooleanSupplier;
+import org.apache.hadoop.fs.CommonConfigurationKeysPublic;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hdds.client.RatisReplicationConfig;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.conf.StorageUnit;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
+import org.apache.hadoop.hdds.scm.pipeline.PipelineManagerImpl;
+import org.apache.hadoop.hdds.utils.IOUtils;
+import org.apache.hadoop.ozone.ClientConfigForTesting;
+import org.apache.hadoop.ozone.DataTestUtil;
+import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.io.SelectorOutputStream;
+import org.apache.hadoop.ozone.om.helpers.BucketLayout;
+import org.apache.ozone.test.GenericTestUtils;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+/**
+ * End-to-end tests for enabling Ratis DataStream on a running cluster
+ * (HDDS-12991). The two tests are isolated (separate clusters):
+ * <ul>
+ * <li>{@link #testDataStreamFallbackAndPortRefresh()}: a streaming client
+ * gracefully falls back to the non-streaming path while a pipeline lacks the
+ * RATIS_DATASTREAM port, and SCM refreshes a datanode's ports when it
+ * re-registers with datastream enabled.</li>
+ * <li>{@link #testCloseNonStreamablePipelineThenStream()}: SCM closes a
+ * pipeline created before datastream (which can never stream in place) so a
+ * fresh streaming-capable pipeline replaces it, after which a streaming
write
+ * succeeds end-to-end.</li>
+ * </ul>
+ *
+ * <p>Writes to portless pipelines throw instead of falling back
+ * (HDDS-12991 part 1 is not yet implemented), so the pre-enable writes in each
+ * test use a non-streaming FileSystem. The post-enable writes use a retry loop
+ * to absorb the transition window while the background scrubber closes
portless
+ * pipelines and a fresh streaming-capable pipeline is created.
+ */
+public class TestOzoneFileSystemDataStreamEnablement {
+
+ // Small threshold/payload keep the writes fast while still selecting the
+ // streaming path (payload > threshold).
+ private static final int AUTO_THRESHOLD = 4 << 10;
+ private static final int WRITE_SIZE = 256 << 10;
+ // Retry budget for writes in the pipeline-transition window.
+ private static final int MAX_WRITE_ATTEMPTS = 10;
+ private static final long WRITE_RETRY_DELAY_MS = 3_000L;
+
+ private MiniOzoneCluster cluster;
+ private OzoneClient client;
+ private OzoneBucket bucket;
+ private OzoneConfiguration conf;
+
+ private void startClusterWithDatanodeStreamDisabled() throws Exception {
+ conf = new OzoneConfiguration();
+ // Datanode side: datastream initially disabled, so pipelines are created
+ // without the RATIS_DATASTREAM port.
+ conf.setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, false);
+ // Client side: always attempt streaming writes.
+ conf.setBoolean(OZONE_FS_DATASTREAM_ENABLED, true);
+ conf.set(OZONE_FS_DATASTREAM_AUTO_THRESHOLD, AUTO_THRESHOLD + "B");
+ conf.setInt(OZONE_SCM_RATIS_PIPELINE_LIMIT, 10);
+ // A long stale interval keeps the OPEN pipeline alive across the (no
+ // stop-wait) rolling restart, so the test drives pipeline closure itself.
+ conf.set(OZONE_SCM_STALENODE_INTERVAL, "5m");
+ conf.set(OZONE_SCM_DEADNODE_INTERVAL, "10m");
+ conf.set(HDDS_HEARTBEAT_INTERVAL, "1s");
+ conf.set(OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL, "1s");
+ // Recreate pipelines quickly after a close so the test does not wait the
+ // default two minutes for a fresh RATIS/THREE pipeline.
+ conf.set(OZONE_SCM_PIPELINE_CREATION_INTERVAL, "1s");
+ // Run the port-scrubber frequently so portless pipelines are closed
quickly.
+ conf.set(OZONE_SCM_PIPELINE_SCRUB_INTERVAL, "5s");
+
+ final int chunkSize = 16 << 10;
+ ClientConfigForTesting.newBuilder(StorageUnit.BYTES)
+ .setChunkSize(chunkSize)
+ .setStreamBufferFlushSize(32 << 10)
+ .setStreamBufferMaxSize(64 << 10)
+ .setDataStreamBufferFlushSize(64 << 10)
+ .setDataStreamMinPacketSize(chunkSize)
+ .setDataStreamWindowSize(5 * chunkSize)
+ .setBlockSize(1 << 20)
+ .applyTo(conf);
+
+ cluster = MiniOzoneCluster.newBuilder(conf).setNumDatanodes(3).build();
+ cluster.waitForClusterToBeReady();
+ client = cluster.newClient();
+ bucket = DataTestUtil.createVolumeAndBucket(client,
+ BucketLayout.FILE_SYSTEM_OPTIMIZED);
+ }
+
+ @AfterEach
+ public void teardown() {
+ IOUtils.closeQuietly(client);
+ if (cluster != null) {
+ cluster.shutdown();
+ }
+ }
+
+ private FileSystem fs() throws IOException {
+ final String rootPath = String.format("%s://%s.%s/",
+ OZONE_URI_SCHEME, bucket.getName(), bucket.getVolumeName());
+ conf.set(CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY, rootPath);
+ return FileSystem.get(conf);
+ }
+
+ /** A FileSystem backed by the same bucket but with datastream disabled. */
+ private FileSystem nonStreamingFs() throws IOException {
+ final String rootPath = String.format("%s://%s.%s/",
+ OZONE_URI_SCHEME, bucket.getName(), bucket.getVolumeName());
+ final OzoneConfiguration noStream = new OzoneConfiguration(conf);
+ noStream.setBoolean(OZONE_FS_DATASTREAM_ENABLED, false);
+ noStream.set(CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY, rootPath);
+ // newInstance bypasses the FileSystem cache so the streaming=false setting
+ // is not shadowed by the streaming=true FS cached under the same URI.
+ return FileSystem.newInstance(noStream);
+ }
+
+ /**
+ * Retries {@link #writeAndGetUnderlying} on IOException to absorb the window
+ * while SCM closes a portless pipeline and a new streaming one is created.
+ */
+ private Class<?> writeWithRetry(FileSystem fs, Path path, byte[] data)
+ throws Exception {
+ for (int attempt = 1; attempt < MAX_WRITE_ATTEMPTS; attempt++) {
+ try {
+ return writeAndGetUnderlying(fs, path, data);
+ } catch (IOException ignored) {
+ Thread.sleep(WRITE_RETRY_DELAY_MS);
+ }
+ }
+ return writeAndGetUnderlying(fs, path, data);
+ }
+
+ /** Write {@code data} and return the underlying stream selected by the FS.
*/
+ private static Class<?> writeAndGetUnderlying(FileSystem fs, Path path,
+ byte[] data) throws IOException {
+ final FSDataOutputStream out = fs.create(path, true);
+ out.write(data);
+ final SelectorOutputStream<?> selector =
+ (SelectorOutputStream<?>) out.getWrappedStream();
+ out.close();
+ return selector.getUnderlying().getClass();
+ }
+
+ private static void assertRoundTrips(FileSystem fs, Path path, byte[]
expected)
+ throws IOException {
+ final byte[] read = new byte[expected.length];
+ try (FSDataInputStream in = fs.open(path)) {
+ in.readFully(read);
+ }
+ assertArrayEquals(expected, read);
+ }
+
+ private static byte[] randomBytes() {
+ final byte[] bytes = new byte[WRITE_SIZE];
+ ThreadLocalRandom.current().nextBytes(bytes);
+ return bytes;
+ }
+
+ private List<Pipeline> openRatisThreePipelines() {
+ return cluster.getStorageContainerManager().getPipelineManager()
+ .getPipelines(RatisReplicationConfig.getInstance(THREE),
+ Pipeline.PipelineState.OPEN);
+ }
+
+ private static boolean allNodesHaveDatastreamPort(Pipeline pipeline) {
+ return pipeline.getNodes().stream()
+ .allMatch(n -> n.hasPort(RATIS_DATASTREAM));
+ }
+
+ /**
+ * Enable datastream on every datanode via a rolling restart. {@code false}
+ * (no stop-wait) keeps each restart short; combined with the long stale
+ * interval the OPEN pipeline survives, so its node snapshot stays portless.
+ * Also reflects the enablement in the SCM config so that
+ * {@code closePipelinesMissingDataStreamPort} does not skip portless
detection.
+ */
+ private void rollingRestartEnablingDataStream() throws Exception {
+ for (int i = 0; i < cluster.getHddsDatanodes().size(); i++) {
+ cluster.getHddsDatanodes().get(i).getConf()
+ .setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ cluster.restartHddsDatanode(i, false);
+ }
+ cluster.waitForClusterToBeReady();
+ // Update the shared conf (picked up by a restarted SCM) and the currently
+ // running SCM so closePipelinesMissingDataStreamPort sees the feature
enabled.
+ conf.setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ cluster.getStorageContainerManager().getConfiguration()
+ .setBoolean(HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true);
+ }
+
+ /** Poll until SCM's node records all expose RATIS_DATASTREAM (validates D).
*/
+ private void waitForAllRegisteredNodesToHaveDatastreamPort()
+ throws InterruptedException, TimeoutException {
+ final BooleanSupplier ready = () -> {
+ final List<? extends DatanodeDetails> nodes = cluster
+ .getStorageContainerManager().getScmNodeManager().getAllNodes();
+ return nodes.size() == cluster.getHddsDatanodes().size()
+ && nodes.stream().allMatch(n -> n.hasPort(RATIS_DATASTREAM));
+ };
+ GenericTestUtils.waitFor(ready, 500, 30_000);
+ }
+
+ /** Poll until an OPEN RATIS/THREE pipeline exposes RATIS_DATASTREAM ports.
*/
+ private void waitForStreamablePipeline()
+ throws InterruptedException, TimeoutException {
+ final BooleanSupplier ready = () -> openRatisThreePipelines().stream()
+ .anyMatch(TestOzoneFileSystemDataStreamEnablement
+ ::allNodesHaveDatastreamPort);
+ GenericTestUtils.waitFor(ready, 500, 60_000);
+ }
+
+ /**
+ * After a rolling restart enables datastream, SCM refreshes the datanodes'
+ * ports. Closing the pre-existing portless pipeline (and replacing it with
+ * a streaming-capable one) may require a few retries because the SCM Ratis
+ * group can briefly lose leadership stability right after the rolling
+ * restart. A write that retries during the transition window eventually
+ * succeeds over the new streaming pipeline.
+ */
+ @Test
+ @Timeout(value = 240, unit = TimeUnit.SECONDS)
+ public void testDataStreamFallbackAndPortRefresh() throws Exception {
+ startClusterWithDatanodeStreamDisabled();
+
+ try (FileSystem fs = fs()) {
+ final byte[] data = randomBytes();
+
+ // Write before enabling datastream. The streaming path would throw on a
+ // portless pipeline (HDDS-12991 part 1 not yet implemented), so use a
+ // non-streaming FS to populate the cluster and open a portless pipeline.
+ final Path before = new Path("/before-enable.dat");
+ try (FileSystem noStream = nonStreamingFs()) {
+ try (FSDataOutputStream out = noStream.create(before, true)) {
+ out.write(data);
+ }
+ }
+ assertRoundTrips(fs, before, data);
+
+ rollingRestartEnablingDataStream();
+ waitForAllRegisteredNodesToHaveDatastreamPort();
+
+ // Close portless pipelines via explicit scrub calls (retried so any
+ // transient SCM Ratis leader disruption from the rolling restart is
+ // absorbed) and wait for a streaming-capable replacement to appear.
+ final PipelineManagerImpl pipelineManager =
+ (PipelineManagerImpl)
cluster.getStorageContainerManager().getPipelineManager();
+ final BooleanSupplier streamingPipelineReady = () -> {
+ pipelineManager.scrubAndClosePipelinesMissingDataStreamPort();
+ return openRatisThreePipelines().stream()
+
.anyMatch(TestOzoneFileSystemDataStreamEnablement::allNodesHaveDatastreamPort);
+ };
+ GenericTestUtils.waitFor(streamingPipelineReady, 1_000, 120_000);
+
+ final Path after = new Path("/after-enable.dat");
+ assertEquals(CapableOzoneFSDataStreamOutput.class,
+ writeWithRetry(fs, after, data));
+ assertRoundTrips(fs, after, data);
+ }
+ }
+
+ /**
+ * A pipeline created while datastream was disabled keeps a portless node
+ * snapshot and a stale datastream address in its Raft group, so it can
+ * never serve streaming even after the datanodes restart. SCM closes it so a
+ * fresh, streaming-capable pipeline is created; a streaming write then
+ * succeeds over the new pipeline (HDDS-12991).
+ */
+ @Test
+ @Timeout(value = 70, unit = TimeUnit.SECONDS)
+ public void testCloseNonStreamablePipelineThenStream() throws Exception {
+ startClusterWithDatanodeStreamDisabled();
+
+ try (FileSystem fs = fs()) {
+ final byte[] data = randomBytes();
+
+ // Create a portless OPEN pipeline. The streaming path would throw on a
+ // portless pipeline (HDDS-12991 part 1 not yet implemented), so use a
+ // non-streaming FS to open a pipeline without the RATIS_DATASTREAM port.
+ final Path p1 = new Path("/legacy.dat");
+ try (FileSystem noStream = nonStreamingFs()) {
+ try (FSDataOutputStream out = noStream.create(p1, true)) {
+ out.write(data);
+ }
+ }
+ final List<Pipeline> before = openRatisThreePipelines();
+ assertFalse(before.isEmpty());
+ before.forEach(p -> assertFalse(allNodesHaveDatastreamPort(p),
+ "pipeline should be portless before enabling datastream"));
+
+ rollingRestartEnablingDataStream();
+ waitForAllRegisteredNodesToHaveDatastreamPort();
+
+ // SCM restart reloads the persisted (still portless) pipeline while the
+ // datanodes are registered with the port (no re-registration event
fires).
+ cluster.restartStorageContainerManager(true);
+ waitForAllRegisteredNodesToHaveDatastreamPort();
+
+ final PipelineManagerImpl pipelineManager =
+ (PipelineManagerImpl)
cluster.getStorageContainerManager().getPipelineManager();
+ final List<Pipeline> reloaded = openRatisThreePipelines();
+ assertFalse(reloaded.isEmpty());
+ reloaded.forEach(p -> assertFalse(allNodesHaveDatastreamPort(p),
+ "reloaded pipeline should still be portless"));
+
+ // Close the pipeline(s) exposing the new datastream port; a fresh
+ // streaming-capable pipeline is created in their place by
+ // BackgroundPipelineCreator.
+ pipelineManager.scrubAndClosePipelinesMissingDataStreamPort();
+ waitForStreamablePipeline();
+
+ // The new pipeline serves a streaming write end-to-end.
+ final Path p2 = new Path("/after-recreate.dat");
+ assertEquals(CapableOzoneFSDataStreamOutput.class,
+ writeWithRetry(fs, p2, data));
+ assertRoundTrips(fs, p2, data);
+ }
+ }
+
+ /**
+ * Full lifecycle over a batch of files: write several files while datastream
+ * is disabled (using a non-streaming FS since the streaming path would throw
+ * on portless pipelines), then enable datastream (rolling restart + SCM
+ * restart + close the non-streamable pipeline), then write several more
files
+ * that must all succeed over a streaming-capable pipeline. Asserts that none
+ * of the writes fail, the post-enablement writes take the streaming path,
and
+ * an OPEN pipeline exposing the RATIS_DATASTREAM port serves them.
+ */
+ @Test
+ @Timeout(value = 120, unit = TimeUnit.SECONDS)
+ public void testBatchWritesAcrossStreamingEnablement() throws Exception {
+ startClusterWithDatanodeStreamDisabled();
+
+ final int fileCount = 5;
+ try (FileSystem fs = fs()) {
+ // Phase 1: datastream disabled. The streaming path throws on portless
+ // pipelines (HDDS-12991 part 1 not yet implemented), so write via a
+ // non-streaming FS to confirm the cluster accepts writes.
+ try (FileSystem noStream = nonStreamingFs()) {
+ for (int i = 0; i < fileCount; i++) {
+ final byte[] data = randomBytes();
+ final Path p = new Path("/disabled-" + i + ".dat");
+ try (FSDataOutputStream out = noStream.create(p, true)) {
+ out.write(data);
+ }
+ assertRoundTrips(noStream, p, data);
+ }
+ }
+
+ // Enable datastream on the datanodes and replace the legacy pipeline.
+ rollingRestartEnablingDataStream();
+ waitForAllRegisteredNodesToHaveDatastreamPort();
+ cluster.restartStorageContainerManager(true);
+ waitForAllRegisteredNodesToHaveDatastreamPort();
+ ((PipelineManagerImpl)
cluster.getStorageContainerManager().getPipelineManager())
+ .scrubAndClosePipelinesMissingDataStreamPort();
+ waitForStreamablePipeline();
+
+ // Phase 2: datastream enabled -> every write streams, none fail.
+ for (int i = 0; i < fileCount; i++) {
+ final byte[] data = randomBytes();
+ final Path p = new Path("/enabled-" + i + ".dat");
+ assertEquals(CapableOzoneFSDataStreamOutput.class,
+ writeWithRetry(fs, p, data),
+ "write after enabling datastream must use the streaming path");
+ assertRoundTrips(fs, p, data);
+ }
+
+ // The post-enablement writes are served by a streaming-capable pipeline.
+ assertTrue(openRatisThreePipelines().stream()
+ .anyMatch(TestOzoneFileSystemDataStreamEnablement
+ ::allNodesHaveDatastreamPort),
+ "an OPEN pipeline should expose the RATIS_DATASTREAM port");
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]