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]

Reply via email to