This is an automated email from the ASF dual-hosted git repository.
errose28 pushed a commit to branch HDDS-14496-zdu
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/HDDS-14496-zdu by this push:
new 44b89fdac5c HDDS-15718. Pass write pipeline version from client to
Datanode on write (#10991)
44b89fdac5c is described below
commit 44b89fdac5c7dd41f5ca82c91026b5765d3e0293
Author: Zita Dombi <[email protected]>
AuthorDate: Wed Sep 9 18:32:52 2026 +0200
HDDS-15718. Pass write pipeline version from client to Datanode on write
(#10991)
---
.../hdds/scm/storage/BlockDataStreamOutput.java | 3 +-
.../hadoop/hdds/protocol/DatanodeDetails.java | 34 +++-
.../apache/hadoop/hdds/scm/pipeline/Pipeline.java | 15 ++
.../hdds/scm/storage/ContainerProtocolCalls.java | 12 +-
.../hadoop/hdds/scm/utils/ClientCommandsUtils.java | 43 +++++
.../scm/storage/TestContainerProtocolCalls.java | 186 +++++++++++++++++++++
.../hdds/scm/utils/TestClientCommandsUtils.java | 77 +++++++++
.../server/ratis/ContainerStateMachine.java | 4 +
.../transport/server/ratis/DispatcherContext.java | 26 +++
.../ozone/container/keyvalue/KeyValueHandler.java | 9 +-
.../server/ratis/TestDispatcherContext.java | 51 ++++++
.../container/keyvalue/TestKeyValueHandler.java | 84 ++++++++++
.../src/main/proto/DatanodeClientProtocol.proto | 3 +
13 files changed, 538 insertions(+), 9 deletions(-)
diff --git
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
index e3fff7528d4..cdf353aff3f 100644
---
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
+++
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
@@ -218,7 +218,8 @@ private DataStreamOutput setupStream(Pipeline pipeline)
throws IOException {
ContainerProtos.ContainerCommandRequestProto.newBuilder()
.setCmdType(streamInitType)
.setContainerID(blockID.get().getContainerID())
- .setDatanodeUuid(id).setWriteChunk(writeChunkRequest);
+ .setDatanodeUuid(id).setWriteChunk(writeChunkRequest)
+ .setWritePipelineVersion(pipeline.getWriteVersion());
if (tokenString != null) {
builder.setEncodedToken(tokenString);
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 c354f7f4d93..befb0ed373b 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
@@ -91,6 +91,10 @@ public class DatanodeDetails extends NodeImpl implements
Comparable<DatanodeDeta
private volatile long persistedOpStateExpiryEpochSec;
private HDDSVersion initialVersion;
private volatile HDDSVersion currentVersion;
+ // The raw serialized currentVersion as it arrived off the wire. Kept
separate from
+ // currentVersion so that a reader which deserializes an unknown newer
version to
+ // UNKNOWN_VERSION still forwards the original value to datanodes instead of
-1.
+ private volatile int serializedCurrentVersion;
private DatanodeDetails(Builder b) {
super(b.hostName, b.networkLocation, NetConstants.NODE_COST_DEFAULT);
@@ -107,6 +111,7 @@ private DatanodeDetails(Builder b) {
persistedOpStateExpiryEpochSec = b.persistedOpStateExpiryEpochSec;
initialVersion = b.initialVersion;
currentVersion = b.currentVersion;
+ serializedCurrentVersion = b.serializedCurrentVersion;
if (b.networkName != null) {
setNetworkName(b.networkName);
}
@@ -135,6 +140,7 @@ public DatanodeDetails(DatanodeDetails datanodeDetails) {
datanodeDetails.getPersistedOpStateExpiryEpochSec();
this.initialVersion = datanodeDetails.getInitialVersion();
this.currentVersion = datanodeDetails.getCurrentVersion();
+ this.serializedCurrentVersion =
datanodeDetails.getSerializedCurrentVersion();
}
public static Codec<DatanodeDetails> getCodec() {
@@ -476,7 +482,7 @@ public static DatanodeDetails.Builder newBuilder(
datanodeDetailsProto.getPersistedOpStateExpiry());
}
if (datanodeDetailsProto.hasCurrentVersion()) {
-
builder.setCurrentVersion(HDDSVersion.deserialize(datanodeDetailsProto.getCurrentVersion()));
+ builder.setCurrentVersion(datanodeDetailsProto.getCurrentVersion());
}
return builder;
}
@@ -613,7 +619,7 @@ public HddsProtos.DatanodeDetailsProto.Builder
toProtoBuilder(
}
}
- builder.setCurrentVersion(versionOverride != null ?
versionOverride.serialize() : currentVersion.serialize());
+ builder.setCurrentVersion(versionOverride != null ?
versionOverride.serialize() : serializedCurrentVersion);
return builder;
}
@@ -660,8 +666,19 @@ public HDDSVersion getCurrentVersion() {
return currentVersion;
}
+ /**
+ * @return the raw serialized currentVersion as it arrived off the wire.
Unlike
+ * {@link #getCurrentVersion()} this preserves the original value even
when it cannot be
+ * deserialized by this component. Use this when forwarding the version
to another
+ * component, e.g. the write pipeline version sent to datanodes.
+ */
+ public int getSerializedCurrentVersion() {
+ return serializedCurrentVersion;
+ }
+
public void setCurrentVersion(HDDSVersion currentVersion) {
this.currentVersion = currentVersion;
+ this.serializedCurrentVersion = currentVersion.serialize();
}
@Override
@@ -746,6 +763,7 @@ public static final class Builder {
private long persistedOpStateExpiryEpochSec = 0;
private HDDSVersion initialVersion = HDDSVersion.DEFAULT_VERSION;
private HDDSVersion currentVersion = HDDSVersion.DEFAULT_VERSION;
+ private int serializedCurrentVersion =
HDDSVersion.DEFAULT_VERSION.serialize();
/**
* Default private constructor. To create Builder instance use
@@ -969,6 +987,18 @@ public Builder setInitialVersion(HDDSVersion v) {
public Builder setCurrentVersion(HDDSVersion v) {
this.currentVersion = v;
+ this.serializedCurrentVersion = v.serialize();
+ return this;
+ }
+
+ /**
+ * Sets the currentVersion from its raw serialized value, preserving the
original even if it
+ * cannot be deserialized (in which case {@link #getCurrentVersion()} is
+ * {@link HDDSVersion#UNKNOWN_VERSION} while {@link
#getSerializedCurrentVersion()} keeps it).
+ */
+ public Builder setCurrentVersion(int serialized) {
+ this.serializedCurrentVersion = serialized;
+ this.currentVersion = HDDSVersion.deserialize(serialized);
return this;
}
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
index 8c084d7c8eb..189a40de9c4 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
@@ -39,6 +39,7 @@
import org.apache.commons.lang3.builder.EqualsBuilder;
import org.apache.commons.lang3.builder.HashCodeBuilder;
import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicatedReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
@@ -289,6 +290,20 @@ public DatanodeDetails getFirstNode(Set<DatanodeDetails>
excluded)
"All nodes are excluded: Pipeline=%s, excluded=%s", id, excluded));
}
+ /**
+ * Returns the serialized component version this pipeline should execute
writes at. All nodes in
+ * a pipeline are expected to share the same version, so we use the first
node's. The raw
+ * serialized value is returned (see {@link
DatanodeDetails#getSerializedCurrentVersion()}) so a
+ * client that cannot deserialize a newer datanode version still forwards
the original value
+ * rather than {@link HDDSVersion#UNKNOWN_VERSION}.
+ */
+ public int getWriteVersion() throws IOException {
+ if (nodeStatus.isEmpty()) {
+ throw new IOException(String.format("Pipeline=%s is empty", id));
+ }
+ return nodeStatus.keySet().iterator().next().getSerializedCurrentVersion();
+ }
+
public DatanodeDetails getClosestNode() throws IOException {
return getClosestNode(null);
}
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
index d2898c395d9..d030db429ae 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
@@ -365,7 +365,8 @@ public static ContainerProtos.FinalizeBlockResponseProto
finalizeBlock(
ContainerCommandRequestProto.newBuilder().setCmdType(Type.FinalizeBlock)
.setContainerID(blockID.getContainerID())
.setDatanodeUuid(id)
- .setFinalizeBlock(finalizeBlockRequest);
+ .setFinalizeBlock(finalizeBlockRequest)
+
.setWritePipelineVersion(xceiverClient.getPipeline().getWriteVersion());
if (token != null) {
builder.setEncodedToken(token.encodeToUrlString());
}
@@ -396,7 +397,8 @@ public static ContainerCommandRequestProto
getPutBlockRequest(
ContainerCommandRequestProto.newBuilder().setCmdType(Type.PutBlock)
.setContainerID(containerBlockData.getBlockID().getContainerID())
.setDatanodeUuid(id)
- .setPutBlock(createBlockRequest);
+ .setPutBlock(createBlockRequest)
+ .setWritePipelineVersion(pipeline.getWriteVersion());
if (tokenString != null) {
builder.setEncodedToken(tokenString);
}
@@ -536,7 +538,8 @@ public static XceiverClientReply writeChunkAsync(
.setCmdType(Type.WriteChunk)
.setContainerID(blockID.getContainerID())
.setDatanodeUuid(id)
- .setWriteChunk(writeChunkRequest);
+ .setWriteChunk(writeChunkRequest)
+
.setWritePipelineVersion(xceiverClient.getPipeline().getWriteVersion());
if (tokenString != null) {
builder.setEncodedToken(tokenString);
@@ -594,7 +597,8 @@ public static PutSmallFileResponseProto writeSmallFile(
.setCmdType(Type.PutSmallFile)
.setContainerID(blockID.getContainerID())
.setDatanodeUuid(id)
- .setPutSmallFile(putSmallFileRequest);
+ .setPutSmallFile(putSmallFileRequest)
+ .setWritePipelineVersion(client.getPipeline().getWriteVersion());
if (token != null) {
builder.setEncodedToken(token.encodeToUrlString());
}
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
index 4edec24edcb..22cd3f0160b 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
@@ -17,13 +17,18 @@
package org.apache.hadoop.hdds.scm.utils;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
/**
* These methods should be merged with other similar utility classes.
*/
public final class ClientCommandsUtils {
+ private static final Logger LOG =
LoggerFactory.getLogger(ClientCommandsUtils.class);
+
/** Utility classes should not be constructed. **/
private ClientCommandsUtils() {
@@ -46,4 +51,42 @@ public static ContainerProtos.ReadChunkVersion
getReadChunkVersion(
return ContainerProtos.ReadChunkVersion.V0;
}
}
+
+ /**
+ * Returns the write pipeline (component) version the datanode should
execute the write at,
+ * derived from the value the client forwarded from the SCM-provided
pipeline version.
+ * {@link HDDSVersion#ZDU} is the lowest version write-path versioning can
use, so it is the
+ * floor: any version below ZDU (including an absent field from a client
predating zero downtime
+ * upgrade support, or a value this datanode cannot deserialize) is rounded
up to ZDU. ZDU is the
+ * first value that is unambiguous across both the component-version and the
(legacy)
+ * layout-feature domains, so any {@code isAllowed} comparison on a datanode
that has not
+ * finalized ZDU is safe. No write-path versioning feature predates ZDU, so
this loses no
+ * behavior. New clients may still forward a pre-ZDU current version (e.g.
when HDDS is not yet
+ * finalized for ZDU); the datanode rounds it up here so there is no issue.
+ */
+ public static HDDSVersion
getWritePipelineVersion(ContainerProtos.ContainerCommandRequestProto request) {
+ final HDDSVersion defaultWriteVersion = HDDSVersion.ZDU;
+ // If no write version is provided, fall back to the default.
+ if (!request.hasWritePipelineVersion()) {
+ return defaultWriteVersion;
+ }
+
+ // If the write version does not deserialize to any known version, fall
back to the default.
+ int serializedVersion = request.getWritePipelineVersion();
+ HDDSVersion writeVersion = HDDSVersion.deserialize(serializedVersion);
+ if (writeVersion == HDDSVersion.UNKNOWN_VERSION) {
+ // Should not normally happen: the version originates from SCM's view of
the datanodes.
+ LOG.error("Datanode was given an unrecognized write pipeline version {};
using {} instead.",
+ serializedVersion, defaultWriteVersion);
+ return defaultWriteVersion;
+ }
+
+ // ZDU is the first version that supports write path versioning.
+ // A lower version should be replaced with the default.
+ if (!defaultWriteVersion.isSupportedBy(writeVersion)) {
+ return defaultWriteVersion;
+ }
+
+ return writeVersion;
+ }
}
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerProtocolCalls.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerProtocolCalls.java
new file mode 100644
index 00000000000..8d934c24e80
--- /dev/null
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerProtocolCalls.java
@@ -0,0 +1,186 @@
+/*
+ * 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.hdds.scm.storage;
+
+import static
org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
+import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.List;
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.apache.hadoop.hdds.client.BlockID;
+import org.apache.hadoop.hdds.client.RatisReplicationConfig;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.BlockData;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type;
+import org.apache.hadoop.hdds.scm.XceiverClientReply;
+import org.apache.hadoop.hdds.scm.XceiverClientSpi;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
+import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
+import org.apache.hadoop.ozone.common.Checksum;
+import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+/**
+ * Test that {@link ContainerProtocolCalls} stamps the pipeline write version
on
+ * outgoing write requests.
+ */
+public class TestContainerProtocolCalls {
+
+ private static final int VERSION = HDDSVersion.ZDU.serialize();
+
+ private static Pipeline pipelineWithVersion(int version) {
+ List<DatanodeDetails> nodes = Arrays.asList(
+ randomDatanodeDetails(), randomDatanodeDetails(),
randomDatanodeDetails());
+ nodes.forEach(dn ->
dn.setCurrentVersion(HDDSVersion.deserialize(version)));
+ return Pipeline.newBuilder()
+ .setId(PipelineID.randomId())
+ .setState(Pipeline.PipelineState.OPEN)
+ .setReplicationConfig(RatisReplicationConfig.getInstance(THREE))
+ .setNodes(nodes)
+ .build();
+ }
+
+ private static XceiverClientSpi clientWith(Pipeline pipeline) {
+ XceiverClientSpi client = mock(XceiverClientSpi.class);
+ when(client.getPipeline()).thenReturn(pipeline);
+ return client;
+ }
+
+ private static ChunkInfo chunkInfo(byte[] data) throws IOException {
+ return ChunkInfo.newBuilder()
+ .setChunkName("chunk")
+ .setOffset(0)
+ .setLen(data.length)
+ .setChecksumData(new Checksum(ChecksumType.CRC32, 256)
+ .computeChecksum(data).getProtoBufMessage())
+ .build();
+ }
+
+ @Test
+ void putBlockRequestCarriesPipelineWriteVersion() throws IOException {
+ Pipeline pipeline = pipelineWithVersion(VERSION);
+ BlockData blockData = BlockData.newBuilder()
+ .setBlockID(DatanodeBlockID.newBuilder()
+ .setContainerID(1).setLocalID(1).build())
+ .build();
+
+ ContainerCommandRequestProto request =
+ ContainerProtocolCalls.getPutBlockRequest(pipeline, blockData, true,
null);
+
+ assertEquals(VERSION, request.getWritePipelineVersion());
+ }
+
+ @Test
+ void writeChunkRequestCarriesPipelineWriteVersion() throws Exception {
+ Pipeline pipeline = pipelineWithVersion(VERSION);
+ XceiverClientSpi client = clientWith(pipeline);
+
when(client.sendCommandAsync(any())).thenReturn(mock(XceiverClientReply.class));
+ byte[] data = "data".getBytes(StandardCharsets.UTF_8);
+
+ ContainerProtocolCalls.writeChunkAsync(client, chunkInfo(data), new
BlockID(1, 1),
+ ByteString.copyFrom(data), null, 0, null, false, true);
+
+ ArgumentCaptor<ContainerCommandRequestProto> captor =
+ ArgumentCaptor.forClass(ContainerCommandRequestProto.class);
+ verify(client).sendCommandAsync(captor.capture());
+ assertEquals(VERSION, captor.getValue().getWritePipelineVersion());
+ }
+
+ @Test
+ void writeSmallFileRequestCarriesPipelineWriteVersion() throws IOException {
+ Pipeline pipeline = pipelineWithVersion(VERSION);
+ XceiverClientSpi client = clientWith(pipeline);
+ when(client.sendCommand(any(), anyList())).thenReturn(
+ response(Type.PutSmallFile));
+
+ ContainerProtocolCalls.writeSmallFile(client, new BlockID(1, 1),
+ "data".getBytes(StandardCharsets.UTF_8), null);
+
+ assertEquals(VERSION, sentRequest(client).getWritePipelineVersion());
+ }
+
+ @Test
+ void finalizeBlockRequestCarriesPipelineWriteVersion() throws IOException {
+ Pipeline pipeline = pipelineWithVersion(VERSION);
+ XceiverClientSpi client = clientWith(pipeline);
+ when(client.sendCommand(any(), anyList())).thenReturn(
+ response(Type.FinalizeBlock));
+
+ ContainerProtocolCalls.finalizeBlock(client,
+ DatanodeBlockID.newBuilder().setContainerID(1).setLocalID(1).build(),
null);
+
+ assertEquals(VERSION, sentRequest(client).getWritePipelineVersion());
+ }
+
+ @Test
+ void writeVersionPreservesOriginalSerializedValueFromWire() throws
IOException {
+ // A datanode current version this (older) client cannot deserialize: it
becomes
+ // UNKNOWN_VERSION locally, but the original serialized int must still be
forwarded.
+ DatanodeDetails dn = DatanodeDetails.getFromProtoBuf(
+ randomDatanodeDetails().getProtoBufMessage().toBuilder()
+ .setCurrentVersion(Integer.MAX_VALUE).build());
+ assertEquals(HDDSVersion.UNKNOWN_VERSION, dn.getCurrentVersion());
+
+ Pipeline pipeline = Pipeline.newBuilder()
+ .setId(PipelineID.randomId())
+ .setState(Pipeline.PipelineState.OPEN)
+ .setReplicationConfig(RatisReplicationConfig.getInstance(THREE))
+ .setNodes(Arrays.asList(dn, randomDatanodeDetails(),
randomDatanodeDetails()))
+ .build();
+ BlockData blockData = BlockData.newBuilder()
+ .setBlockID(DatanodeBlockID.newBuilder()
+ .setContainerID(1).setLocalID(1).build())
+ .build();
+
+ ContainerCommandRequestProto request =
+ ContainerProtocolCalls.getPutBlockRequest(pipeline, blockData, true,
null);
+
+ assertEquals(Integer.MAX_VALUE, request.getWritePipelineVersion());
+ }
+
+ private static ContainerCommandResponseProto response(Type type) {
+ return ContainerCommandResponseProto.newBuilder()
+ .setCmdType(type)
+ .setResult(Result.SUCCESS)
+ .build();
+ }
+
+ private static ContainerCommandRequestProto sentRequest(XceiverClientSpi
client)
+ throws IOException {
+ ArgumentCaptor<ContainerCommandRequestProto> captor =
+ ArgumentCaptor.forClass(ContainerCommandRequestProto.class);
+ verify(client).sendCommand(captor.capture(), anyList());
+ return captor.getValue();
+ }
+}
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/utils/TestClientCommandsUtils.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/utils/TestClientCommandsUtils.java
new file mode 100644
index 00000000000..b37c86085bc
--- /dev/null
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/utils/TestClientCommandsUtils.java
@@ -0,0 +1,77 @@
+/*
+ * 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.hdds.scm.utils;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import org.apache.hadoop.hdds.HDDSVersion;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test for {@link ClientCommandsUtils}.
+ */
+public class TestClientCommandsUtils {
+
+ private static ContainerCommandRequestProto.Builder baseRequest() {
+ return ContainerCommandRequestProto.newBuilder()
+ .setCmdType(Type.WriteChunk)
+ .setContainerID(1)
+ .setDatanodeUuid("dn");
+ }
+
+ @Test
+ void returnsWritePipelineVersionWhenAtLeastZdu() {
+ ContainerCommandRequestProto request = baseRequest()
+ .setWritePipelineVersion(HDDSVersion.ZDU.serialize())
+ .build();
+
+ assertEquals(HDDSVersion.ZDU,
+ ClientCommandsUtils.getWritePipelineVersion(request));
+ }
+
+ @Test
+ void roundsUpToZduWhenBelowZdu() {
+ // A recognized pre-ZDU version must be rounded up to ZDU (the
write-versioning floor).
+ ContainerCommandRequestProto request = baseRequest()
+ .setWritePipelineVersion(HDDSVersion.SHORT_CIRCUIT_READS.serialize())
+ .build();
+
+ assertEquals(HDDSVersion.ZDU,
+ ClientCommandsUtils.getWritePipelineVersion(request));
+ }
+
+ @Test
+ void defaultsToZduWhenAbsent() {
+ ContainerCommandRequestProto request = baseRequest().build();
+
+ assertEquals(HDDSVersion.ZDU,
+ ClientCommandsUtils.getWritePipelineVersion(request));
+ }
+
+ @Test
+ void fallsBackToZduForUnknownVersion() {
+ ContainerCommandRequestProto request = baseRequest()
+ .setWritePipelineVersion(HDDSVersion.UNKNOWN_VERSION.serialize())
+ .build();
+
+ assertEquals(HDDSVersion.ZDU,
+ ClientCommandsUtils.getWritePipelineVersion(request));
+ }
+}
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
index 9a3e988e27a..5db6985b1f2 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
@@ -68,6 +68,7 @@
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
import
org.apache.hadoop.hdds.scm.container.common.helpers.ContainerNotOpenException;
import
org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
+import org.apache.hadoop.hdds.scm.utils.ClientCommandsUtils;
import org.apache.hadoop.hdds.utils.Cache;
import org.apache.hadoop.hdds.utils.ResourceCache;
import org.apache.hadoop.ozone.HddsDatanodeService;
@@ -605,6 +606,7 @@ private CompletableFuture<Message> writeStateMachineData(
.setLogIndex(entryIndex)
.setStage(DispatcherContext.WriteChunkStage.WRITE_DATA)
.setContainer2BCSIDMap(container2BCSIDMap)
+
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(requestProto))
.build();
CompletableFuture<Message> raftFuture = new CompletableFuture<>();
// ensure the write chunk happens asynchronously in writeChunkExecutor
pool thread.
@@ -751,6 +753,7 @@ public CompletableFuture<DataStream>
stream(RaftClientRequest request) {
.newBuilder(DispatcherContext.Op.STREAM_INIT)
.setStage(DispatcherContext.WriteChunkStage.WRITE_DATA)
.setContainer2BCSIDMap(container2BCSIDMap)
+
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(requestProto))
.build();
final DataChannel channel = getStreamDataChannel(requestProto,
context);
final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ?
@@ -1150,6 +1153,7 @@ public CompletableFuture<Message>
applyTransaction(TransactionContext trx) {
Objects.requireNonNull(context, "context == null");
final ContainerCommandRequestProto requestProto = context.getLogProto();
final Type cmdType = requestProto.getCmdType();
+
builder.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(requestProto));
// Make sure that in write chunk, the user data is not set
if (cmdType == Type.WriteChunk) {
Preconditions
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
index e0f0ccbe193..0bf3221fde0 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
@@ -19,6 +19,7 @@
import java.util.Map;
import java.util.Objects;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.annotation.InterfaceAudience;
import org.apache.hadoop.hdds.annotation.InterfaceStability;
import org.apache.hadoop.util.Time;
@@ -53,6 +54,10 @@ public final class DispatcherContext {
private final Map<Long, Long> container2BCSIDMap;
+ // the component version the client asked the datanode to execute this write
at
+ // defaults to HDDSVersion.ZDU (the safe cross-domain floor) when the client
did not supply one
+ private final HDDSVersion writeVersion;
+
private final boolean releaseSupported;
private volatile Runnable releaseMethod;
@@ -137,6 +142,7 @@ private DispatcherContext(Builder b) {
this.logIndex = b.logIndex;
this.stage = b.stage;
this.container2BCSIDMap = b.container2BCSIDMap;
+ this.writeVersion = b.writeVersion;
this.releaseSupported = b.releaseSupported;
}
@@ -161,6 +167,14 @@ public Map<Long, Long> getContainer2BCSIDMap() {
return container2BCSIDMap;
}
+ /**
+ * @return the component version the client requested this write be executed
at.
+ * Defaults to {@link HDDSVersion#ZDU}, the safe cross-domain floor,
when the client did not supply one.
+ */
+ public HDDSVersion getWriteVersion() {
+ return writeVersion;
+ }
+
public long getStartTime() {
return startTime;
}
@@ -198,6 +212,7 @@ public static final class Builder {
private long term;
private long logIndex;
private Map<Long, Long> container2BCSIDMap;
+ private HDDSVersion writeVersion = HDDSVersion.ZDU;
private boolean releaseSupported;
private Builder(Op op) {
@@ -248,6 +263,17 @@ public Builder setContainer2BCSIDMap(Map<Long, Long> map) {
return this;
}
+ /**
+ * Sets the component version the client requested this write be executed
at.
+ *
+ * @param version the write pipeline (component) version
+ * @return Builder
+ */
+ public Builder setWriteVersion(HDDSVersion version) {
+ this.writeVersion = version;
+ return this;
+ }
+
public Builder setReleaseSupported(boolean releaseSupported) {
this.releaseSupported = releaseSupported;
return this;
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
index a32a56a0267..d361b2f5c61 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
@@ -123,6 +123,7 @@
import org.apache.hadoop.hdds.scm.storage.BlockInputStream;
import org.apache.hadoop.hdds.scm.storage.BlockLocationInfo;
import org.apache.hadoop.hdds.scm.storage.ChunkInputStream;
+import org.apache.hadoop.hdds.scm.utils.ClientCommandsUtils;
import org.apache.hadoop.hdds.security.token.OzoneBlockTokenIdentifier;
import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
import org.apache.hadoop.hdds.utils.FaultInjector;
@@ -1106,7 +1107,9 @@ ContainerCommandResponseProto handleWriteChunk(
ChunkBuffer data = null;
if (dispatcherContext == null) {
- dispatcherContext = DispatcherContext.getHandleWriteChunk();
+ dispatcherContext =
DispatcherContext.newBuilder(DispatcherContext.Op.HANDLE_WRITE_CHUNK)
+
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(request))
+ .build();
}
final boolean isWrite = dispatcherContext.getStage().isWrite();
if (isWrite) {
@@ -1265,7 +1268,9 @@ ContainerCommandResponseProto handlePutSmallFile(
ChunkBuffer data = ChunkBuffer.wrap(
putSmallFileReq.getData().asReadOnlyByteBufferList());
if (dispatcherContext == null) {
- dispatcherContext = DispatcherContext.getHandlePutSmallFile();
+ dispatcherContext =
DispatcherContext.newBuilder(DispatcherContext.Op.HANDLE_PUT_SMALL_FILE)
+
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(request))
+ .build();
}
BlockID blockID = blockData.getBlockID();
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/TestDispatcherContext.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/TestDispatcherContext.java
new file mode 100644
index 00000000000..b8edc6de325
--- /dev/null
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/TestDispatcherContext.java
@@ -0,0 +1,51 @@
+/*
+ * 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.ozone.container.common.transport.server.ratis;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test for {@link DispatcherContext} write version handling (HDDS-15718).
+ */
+public class TestDispatcherContext {
+
+ @Test
+ void writeVersionDefaultsToZdu() {
+ DispatcherContext context = DispatcherContext
+ .newBuilder(DispatcherContext.Op.WRITE_STATE_MACHINE_DATA)
+ .build();
+
+ assertEquals(HDDSVersion.ZDU,
+ context.getWriteVersion());
+ }
+
+ @Test
+ void writeVersionIsCarried() {
+ // Use a non-default version so this exercises the setter, not the ZDU
default.
+ DispatcherContext context = DispatcherContext
+ .newBuilder(DispatcherContext.Op.WRITE_STATE_MACHINE_DATA)
+ .setWriteVersion(HDDSVersion.SHORT_CIRCUIT_READS)
+ .build();
+
+ assertEquals(HDDSVersion.SHORT_CIRCUIT_READS,
+ context.getWriteVersion());
+ }
+}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
index 2670f1585a2..92cb3adad12 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
@@ -56,6 +56,7 @@
import com.google.common.collect.Sets;
import java.io.File;
import java.io.IOException;
+import java.lang.reflect.Field;
import java.nio.ByteBuffer;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -72,6 +73,7 @@
import org.apache.commons.io.FileUtils;
import org.apache.hadoop.conf.StorageUnit;
import org.apache.hadoop.fs.FileUtil;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.BlockID;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
@@ -84,6 +86,7 @@
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto;
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
import
org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
+import org.apache.hadoop.hdds.scm.pipeline.MockPipeline;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
import org.apache.hadoop.hdds.security.token.TokenVerifier;
import org.apache.hadoop.hdds.utils.io.RandomAccessFileChannel;
@@ -113,6 +116,7 @@
import
org.apache.hadoop.ozone.container.common.volume.RoundRobinVolumeChoosingPolicy;
import org.apache.hadoop.ozone.container.common.volume.StorageVolume;
import org.apache.hadoop.ozone.container.common.volume.VolumeSet;
+import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager;
import org.apache.hadoop.ozone.container.ozoneimpl.ContainerController;
import
org.apache.hadoop.ozone.container.ozoneimpl.ContainerScannerConfiguration;
import org.apache.hadoop.ozone.container.ozoneimpl.OnDemandContainerScanner;
@@ -124,6 +128,7 @@
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -994,6 +999,85 @@ private HandlerWithVolumeSet createKeyValueHandler(Path
path) throws IOException
return new HandlerWithVolumeSet(kvHandler, volumeSet, containerSet);
}
+ /**
+ * EC/standalone writes reach the datanode with a null DispatcherContext, so
KeyValueHandler
+ * fabricates one. Verify the client-supplied write pipeline version is
carried onto that
+ * context for WriteChunk.
+ */
+ @Test
+ public void testWriteChunkCarriesWritePipelineVersionForEc() throws
Exception {
+ Path testDir = Files.createTempDirectory("testWriteChunkEcVersion");
+ conf.set(OZONE_SCM_CONTAINER_LAYOUT_KEY,
ContainerLayoutVersion.FILE_PER_BLOCK.name());
+ HandlerWithVolumeSet handlerWithVolume =
createKeyValueHandlerWithVolumeSet(testDir);
+ KeyValueHandler kvHandler = handlerWithVolume.getHandler();
+
+ long containerID = ContainerTestHelper.getTestContainerID();
+ KeyValueContainer container = createOpenContainer(containerID,
+ handlerWithVolume.getVolumeSet(), handlerWithVolume.getContainerSet());
+ ChunkManager spyChunkManager = spyChunkManager(kvHandler);
+
+ BlockID blockID = ContainerTestHelper.getTestBlockID(containerID);
+ ContainerCommandRequestProto request = ContainerTestHelper
+ .getWriteChunkRequest(MockPipeline.createSingleNodePipeline(),
blockID, 1024)
+ .toBuilder()
+ .setWritePipelineVersion(HDDSVersion.ZDU.serialize())
+ .build();
+
+ kvHandler.handleWriteChunk(request, container, null);
+
+ ArgumentCaptor<DispatcherContext> captor =
ArgumentCaptor.forClass(DispatcherContext.class);
+ verify(spyChunkManager).writeChunk(eq(container), any(), any(),
any(ChunkBuffer.class), captor.capture());
+ assertEquals(HDDSVersion.ZDU, captor.getValue().getWriteVersion());
+ }
+
+ /**
+ * Same as {@link #testWriteChunkCarriesWritePipelineVersionForEc()} but for
PutSmallFile.
+ */
+ @Test
+ public void testPutSmallFileCarriesWritePipelineVersionForEc() throws
Exception {
+ Path testDir = Files.createTempDirectory("testPutSmallFileEcVersion");
+ conf.set(OZONE_SCM_CONTAINER_LAYOUT_KEY,
ContainerLayoutVersion.FILE_PER_BLOCK.name());
+ HandlerWithVolumeSet handlerWithVolume =
createKeyValueHandlerWithVolumeSet(testDir);
+ KeyValueHandler kvHandler = handlerWithVolume.getHandler();
+
+ long containerID = ContainerTestHelper.getTestContainerID();
+ KeyValueContainer container = createOpenContainer(containerID,
+ handlerWithVolume.getVolumeSet(), handlerWithVolume.getContainerSet());
+ ChunkManager spyChunkManager = spyChunkManager(kvHandler);
+
+ BlockID blockID = ContainerTestHelper.getTestBlockID(containerID);
+ ContainerCommandRequestProto request = ContainerTestHelper
+ .getWriteSmallFileRequest(MockPipeline.createSingleNodePipeline(),
blockID, 1024)
+ .toBuilder()
+ .setWritePipelineVersion(HDDSVersion.ZDU.serialize())
+ .build();
+
+ kvHandler.handlePutSmallFile(request, container, null);
+
+ ArgumentCaptor<DispatcherContext> captor =
ArgumentCaptor.forClass(DispatcherContext.class);
+ verify(spyChunkManager).writeChunk(eq(container), any(), any(),
any(ChunkBuffer.class), captor.capture());
+ assertEquals(HDDSVersion.ZDU, captor.getValue().getWriteVersion());
+ }
+
+ private static ChunkManager spyChunkManager(KeyValueHandler kvHandler)
throws Exception {
+ Field field = KeyValueHandler.class.getDeclaredField("chunkManager");
+ field.setAccessible(true);
+ ChunkManager spy = spy((ChunkManager) field.get(kvHandler));
+ field.set(kvHandler, spy);
+ return spy;
+ }
+
+ private KeyValueContainer createOpenContainer(long containerID,
+ MutableVolumeSet volumeSet, ContainerSet containerSet) throws
IOException {
+ KeyValueContainerData containerData = new KeyValueContainerData(
+ containerID, ContainerLayoutVersion.FILE_PER_BLOCK,
+ (long) StorageUnit.GB.toBytes(1), UUID.randomUUID().toString(),
DATANODE_UUID);
+ KeyValueContainer container = new KeyValueContainer(containerData, conf);
+ container.create(volumeSet, new RoundRobinVolumeChoosingPolicy(),
CLUSTER_ID);
+ containerSet.addContainer(container);
+ return container;
+ }
+
private static class HandlerWithVolumeSet {
private final KeyValueHandler handler;
private final MutableVolumeSet volumeSet;
diff --git
a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
index 72d50c975a5..385cdb319cf 100644
--- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
+++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
@@ -203,6 +203,9 @@ message ContainerCommandRequestProto {
optional GetContainerChecksumInfoRequestProto getContainerChecksumInfo =
27;
optional ReadBlockRequestProto readBlock = 28;
+ // Version that datanodes should use for this write operation.
+ // Set by SCM and forwarded by the client to the datanodes.
+ optional uint32 writePipelineVersion = 29;
// clientId and callId are used to distinguish different requests from
different local clients for shortCircuitRead
optional bytes clientId = 100;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]