From 490277368bfae0ee80b34b5fc69c3ffde29aa288 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Thu, 9 Jul 2026 09:41:02 +0800 Subject: [PATCH 1/8] init --- .../statemachine/DatanodeConfiguration.java | 19 +++ .../server/ratis/ContainerStateMachine.java | 29 +++- .../impl/KeyValueStreamDataChannel.java | 105 +++++++++++++- .../keyvalue/impl/StreamDataChannelBase.java | 4 + .../impl/StreamPutBlockProcessor.java | 29 ++++ .../TestDatanodeConfiguration.java | 17 +++ .../impl/TestKeyValueStreamDataChannel.java | 137 ++++++++++++------ 7 files changed, 289 insertions(+), 51 deletions(-) create mode 100644 hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java index 506dd79c37dc..276cbbd561d8 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java @@ -535,6 +535,17 @@ public class DatanodeConfiguration extends ReconfigurableConfig { private boolean waitOnAllFollowers = WAIT_ON_ALL_FOLLOWERS_DEFAULT; + public static final String HDDS_DATANODE_DATASTREAM_PUTBLOCK_ENABLED = + CONFIG_PREFIX + ".datastream.putblock.enabled"; + + @Config(key = "hdds.datanode.datastream.putblock.enabled", + defaultValue = "false", + type = ConfigType.BOOLEAN, + tags = { DATANODE }, + description = "When enabled, PutBlock is committed when a Ratis data stream " + + "closes instead of via the Raft WriteAsync path.") + private boolean datastreamPutBlockEnabled = false; + @Config(key = "hdds.datanode.container.schema.v3.enabled", defaultValue = "true", type = ConfigType.BOOLEAN, @@ -1310,4 +1321,12 @@ public int getGrpcSoBacklog() { public void setGrpcSoBacklog(int grpcSoBacklog) { this.grpcSoBacklog = grpcSoBacklog; } + + public boolean isDatastreamPutBlockEnabled() { + return datastreamPutBlockEnabled; + } + + public void setDatastreamPutBlockEnabled(boolean datastreamPutBlockEnabled) { + this.datastreamPutBlockEnabled = datastreamPutBlockEnabled; + } } 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 3de4110a01b3..3578df557f61 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 @@ -160,6 +160,7 @@ public class ContainerStateMachine extends BaseStateMachine { private final Semaphore applyTransactionSemaphore; private final boolean waitOnBothFollowers; + private final boolean datastreamPutBlockEnabled; private final HddsDatanodeService datanodeService; private static Semaphore semaphore = new Semaphore(1); private final AtomicBoolean peersValidated; @@ -302,6 +303,9 @@ public ContainerStateMachine(HddsDatanodeService hddsDatanodeService, RaftGroupI this.waitOnBothFollowers = conf.getObject( DatanodeConfiguration.class).waitOnAllFollowers(); + this.datastreamPutBlockEnabled = conf.getObject( + DatanodeConfiguration.class).isDatastreamPutBlockEnabled(); + this.writeChunkWaitMaxNs = conf.getTimeDuration(ScmConfigKeys.HDDS_CONTAINER_RATIS_STATEMACHINE_WRITE_WAIT_INTERVAL, ScmConfigKeys.HDDS_CONTAINER_RATIS_STATEMACHINE_WRITE_WAIT_INTERVAL_NS_DEFAULT, TimeUnit.NANOSECONDS); } @@ -742,6 +746,15 @@ public CompletableFuture stream(RaftClientRequest request) { .setContainer2BCSIDMap(container2BCSIDMap) .build(); DataChannel channel = getStreamDataChannel(requestProto, context); + if (datastreamPutBlockEnabled && channel instanceof KeyValueStreamDataChannel) { + final KeyValueStreamDataChannel kvChannel = (KeyValueStreamDataChannel) channel; + kvChannel.setDatastreamPutBlockEnabled(true); + kvChannel.setPutBlockProcessor(req -> dispatchCommand(req, + DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_LINK) + .setStage(DispatcherContext.WriteChunkStage.COMBINED) + .setContainer2BCSIDMap(container2BCSIDMap) + .build())); + } final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ? getChunkExecutor(requestProto.getWriteChunk()) : null; return new LocalStream(channel, chunkExecutor); @@ -773,8 +786,20 @@ public CompletableFuture link(DataStream stream, LogEntryProto entry) { final KeyValueStreamDataChannel kvStreamDataChannel = (KeyValueStreamDataChannel) dataChannel; - kvStreamDataChannel.setLinked(); - return CompletableFuture.completedFuture(null); + + if (datastreamPutBlockEnabled) { + // The PutBlock should be commited when the stream is closed thus + // we expect this stream has been marked as linked. + if (kvStreamDataChannel.isLinked()) { + return CompletableFuture.completedFuture(null); + } else { + return JavaUtils.completeExceptionally(new IllegalStateException( + "PutBlock was not committed on stream close: " + kvStreamDataChannel)); + } + } else { + kvStreamDataChannel.setLinked(); + return CompletableFuture.completedFuture(null); + } } private ExecutorService getChunkExecutor(WriteChunkRequestProto req) { diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java index 3218fb4f88d9..73027e65ee81 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java @@ -21,13 +21,18 @@ import java.io.File; import java.io.IOException; import java.nio.ByteBuffer; +import java.util.Objects; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; +import org.apache.hadoop.hdds.ratis.ContainerCommandRequestMessage; import org.apache.hadoop.hdds.ratis.RatisHelper; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.hdds.scm.storage.BlockDataStreamOutput; import org.apache.hadoop.ozone.container.common.helpers.ContainerMetrics; import org.apache.hadoop.ozone.container.common.impl.ContainerData; +import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf; import org.apache.ratis.util.ReferenceCountedObject; import org.slf4j.Logger; @@ -41,12 +46,25 @@ public class KeyValueStreamDataChannel extends StreamDataChannelBase { private final Buffers buffers = new Buffers(BlockDataStreamOutput.PUT_BLOCK_REQUEST_LENGTH_MAX); + private final AtomicReference putBlockRequest + = new AtomicReference<>(); private final AtomicBoolean closed = new AtomicBoolean(); + private boolean datastreamPutBlockEnabled; + private StreamPutBlockProcessor putBlockProcessor; KeyValueStreamDataChannel(File file, ContainerData containerData, ContainerMetrics metrics) throws StorageContainerException { super(file, containerData, metrics); + datastreamPutBlockEnabled = false; + } + + public void setDatastreamPutBlockEnabled(boolean enabled) { + this.datastreamPutBlockEnabled = enabled; + } + + public void setPutBlockProcessor(StreamPutBlockProcessor processor) { + this.putBlockProcessor = processor; } @Override @@ -87,6 +105,11 @@ static void writeFully(ByteBuffer b, WriteMethod writeMethod) } } + public ContainerCommandRequestProto getPutBlockRequest() { + return Objects.requireNonNull(putBlockRequest.get(), + () -> "putBlockRequest == null, " + this); + } + void assertOpen() throws IOException { if (closed.get()) { throw new IOException("Already closed: " + this); @@ -97,7 +120,13 @@ void assertOpen() throws IOException { public void close() throws IOException { if (closed.compareAndSet(false, true)) { try { - writeBuffers(); + if (datastreamPutBlockEnabled) { + // This path requires the client appends the PutBlock at the + // end of the stream. + closeWithStreamPutBlock(); + } else { + writeBuffers(); + } } finally { super.close(); } @@ -128,6 +157,21 @@ private void writeBuffers() throws IOException { } } + static ContainerCommandRequestProto closeBuffers( + Buffers buffers, WriteMethod writeMethod) throws IOException { + final ReferenceCountedObject ref = buffers.pollAll(); + final ByteBuf buf = ref.retain(); + final ContainerCommandRequestProto putBlockRequestProto; + try { + putBlockRequestProto = readPutBlockRequest(buf); + // write the remaining data + writeFully(buf.nioBuffer(), writeMethod); + } finally { + ref.release(); + } + return putBlockRequestProto; + } + static int readProtoLength(ByteBuf b, int lengthIndex) { final int readerIndex = b.readerIndex(); LOG.debug("{}, lengthIndex = {}, readerIndex = {}", @@ -155,6 +199,65 @@ static void setEndIndex(ByteBuf b) { b.writerIndex(protoIndex); } + static ContainerCommandRequestProto readPutBlockRequest(ByteBuf b) + throws IOException { + // readerIndex protoIndex lengthIndex readerIndex+readableBytes + // V V V V + // format: |--- data ---|--- proto ---|--- proto length (4 bytes) ---| + final int readerIndex = b.readerIndex(); + final int lengthIndex = readerIndex + b.readableBytes() - 4; + if (lengthIndex < readerIndex) { + throw new IOException("Buffer too short for PutBlock length: " + b.readableBytes()); + } + final int protoLength = readProtoLength(b.duplicate(), lengthIndex); + final int protoIndex = lengthIndex - protoLength; + if (protoIndex < readerIndex) { + throw new IOException("Invalid PutBlock proto length: " + protoLength); + } + + final ContainerCommandRequestProto proto; + try { + proto = readPutBlockRequest(b.slice(protoIndex, protoLength).nioBuffer()); + } catch (Throwable t) { + RatisHelper.debug(b, "catch", LOG); + throw new IOException("Failed to readPutBlockRequest from " + b + + ": readerIndex=" + readerIndex + + ", protoIndex=" + protoIndex + + ", protoLength=" + protoLength + + ", lengthIndex=" + lengthIndex, t); + } + + // set index for reading data + b.writerIndex(protoIndex); + + return proto; + } + + private static ContainerCommandRequestProto readPutBlockRequest(ByteBuffer b) + throws IOException { + RatisHelper.debug(b, "readPutBlockRequest", LOG); + final ByteString byteString = ByteString.copyFrom(b); + + final ContainerCommandRequestProto request = + ContainerCommandRequestMessage.toProto(byteString, null); + + if (!request.hasPutBlock()) { + throw new StorageContainerException( + "Malformed PutBlock request. trace ID: " + request.getTraceID(), + ContainerProtos.Result.MALFORMED_REQUEST); + } + return request; + } + + private void closeWithStreamPutBlock() throws IOException { + final ContainerCommandRequestProto proto = + closeBuffers(buffers, super::writeFileChannel); + putBlockRequest.set(proto); + Preconditions.checkState(putBlockProcessor == null); + putBlockProcessor.processPutBlock(proto); + setLinked(); + } + interface WriteMethod { int applyAsInt(ByteBuffer src) throws IOException; } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java index 43bcea5e9bd9..3a65c3c42f5c 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java @@ -102,6 +102,10 @@ public void setLinked() { linked.set(true); } + public boolean isLinked() { + return linked.get(); + } + /** * @return true if {@link org.apache.ratis.statemachine.StateMachine.DataChannel} is already linked. */ diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java new file mode 100644 index 000000000000..fccf661bceed --- /dev/null +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java @@ -0,0 +1,29 @@ +/* + * 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.keyvalue.impl; + +import java.io.IOException; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; + +/** + * Processes a PutBlock request received inline from the streaming write path. + */ +@FunctionalInterface +public interface StreamPutBlockProcessor { + void processPutBlock(ContainerCommandRequestProto request) throws IOException; +} diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java index 59e90ba30b38..288a99095670 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java @@ -359,4 +359,21 @@ void testGrpcSoBacklogSetter() { assertEquals(512, subject.getGrpcSoBacklog()); } + + @Test + void testDatastreamPutBlockEnabledDefault() { + DatanodeConfiguration subject = new OzoneConfiguration() + .getObject(DatanodeConfiguration.class); + assertThat(subject.isDatastreamPutBlockEnabled()).isFalse(); + } + + @Test + void testDatastreamPutBlockConfigParsing() { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setBoolean(DatanodeConfiguration.HDDS_DATANODE_DATASTREAM_PUTBLOCK_ENABLED, true); + + DatanodeConfiguration subject = conf.getObject(DatanodeConfiguration.class); + + assertThat(subject.isDatastreamPutBlockEnabled()).isTrue(); + } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index 3c2992f36f4b..08263a8c8f15 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -23,7 +23,10 @@ import static org.apache.hadoop.ozone.container.keyvalue.impl.KeyValueStreamDataChannel.writeBuffers; import static org.apache.hadoop.ozone.container.keyvalue.impl.KeyValueStreamDataChannel.writeFully; import static org.assertj.core.api.Assertions.assertThat; +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.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; @@ -33,6 +36,7 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.nio.channels.WritableByteChannel; +import java.nio.file.Files; import java.util.ArrayList; import java.util.Collection; import java.util.List; @@ -41,16 +45,15 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.atomic.AtomicReference; import org.apache.commons.lang3.RandomUtils; import org.apache.hadoop.hdds.fs.SpaceUsageSource; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.BlockData; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.PutBlockRequestProto; -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.ratis.ContainerCommandRequestMessage; -import org.apache.hadoop.hdds.ratis.RatisHelper; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.container.common.helpers.ContainerMetrics; @@ -66,11 +69,12 @@ import org.apache.ratis.protocol.ClientId; import org.apache.ratis.protocol.DataStreamReply; import org.apache.ratis.protocol.RaftClientReply; -import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf; import org.apache.ratis.thirdparty.io.netty.buffer.Unpooled; import org.apache.ratis.util.ReferenceCountedObject; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -95,6 +99,61 @@ public class TestKeyValueStreamDataChannel { LOG.info("PUT_BLOCK_PROTO_SIZE = {}", PUT_BLOCK_PROTO_SIZE); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testClosePutBlockBehavior(boolean datastreamPutBlockEnabled) throws Exception { + File tempFile = File.createTempFile("test-kv-close-" + datastreamPutBlockEnabled, ".tmp"); + tempFile.deleteOnExit(); + KeyValueStreamDataChannel channel = newChannel(tempFile); + AtomicReference processed = new AtomicReference<>(); + if (datastreamPutBlockEnabled) { + channel.setDatastreamPutBlockEnabled(datastreamPutBlockEnabled); + channel.setPutBlockProcessor(processed::set); + } + + final byte[] data = RandomUtils.secure().randomBytes(50); + final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( + PUT_BLOCK_PROTO, null).getContent().asReadOnlyByteBuffer(); + final ByteBuffer protoLengthBuf = + getProtoLength(putBlockBuf, PUT_BLOCK_REQUEST_LENGTH_MAX); + + write(channel, data); + write(channel, putBlockBuf.duplicate()); + write(channel, protoLengthBuf.duplicate()); + + assertThat(processed.get()).isNull(); + channel.close(); + + if (datastreamPutBlockEnabled) { + assertEquals(PUT_BLOCK_PROTO, processed.get()); + assertEquals(PUT_BLOCK_PROTO, channel.getPutBlockRequest()); + assertTrue(channel.isLinked()); + } else { + assertThat(processed.get()).isNull(); + assertFalse(channel.isLinked()); + } + assertEquals(data.length, tempFile.length()); + assertArrayEquals(data, Files.readAllBytes(tempFile.toPath())); + } + + @Test + public void testReadPutBlockRequestBufferTooShort() { + final ByteBuf buf = Unpooled.buffer(2); + buf.writeByte(1); + buf.writeByte(2); + assertThrows(IOException.class, () -> KeyValueStreamDataChannel.readPutBlockRequest(buf)); + buf.release(); + } + + @Test + public void testReadPutBlockRequestInvalidProtoLength() { + final ByteBuf buf = Unpooled.buffer(8); + buf.writeInt(1); + buf.writeInt(100); + assertThrows(IOException.class, () -> KeyValueStreamDataChannel.readPutBlockRequest(buf)); + buf.release(); + } + @Test public void testSerialization() throws Exception { final int max = PUT_BLOCK_REQUEST_LENGTH_MAX; @@ -112,53 +171,10 @@ public void testSerialization() throws Exception { buf.writeBytes(putBlockBuf); buf.writeBytes(protoLengthBuf); - final ContainerCommandRequestProto proto = readPutBlockRequest(buf); + final ContainerCommandRequestProto proto = KeyValueStreamDataChannel.readPutBlockRequest(buf); assertEquals(PUT_BLOCK_PROTO, proto); } - static ContainerCommandRequestProto readPutBlockRequest(ByteBuf b) throws IOException { - // readerIndex protoIndex lengthIndex readerIndex+readableBytes - // V V V V - // format: |--- data ---|--- proto ---|--- proto length (4 bytes) ---| - final int readerIndex = b.readerIndex(); - final int lengthIndex = readerIndex + b.readableBytes() - 4; - final int protoLength = KeyValueStreamDataChannel.readProtoLength(b.duplicate(), lengthIndex); - final int protoIndex = lengthIndex - protoLength; - - final ContainerCommandRequestProto proto; - try { - proto = readPutBlockRequest(b.slice(protoIndex, protoLength).nioBuffer()); - } catch (Throwable t) { - RatisHelper.debug(b, "catch", LOG); - throw new IOException("Failed to readPutBlockRequest from " + b - + ": readerIndex=" + readerIndex - + ", protoIndex=" + protoIndex - + ", protoLength=" + protoLength - + ", lengthIndex=" + lengthIndex, t); - } - - // set index for reading data - b.writerIndex(protoIndex); - - return proto; - } - - private static ContainerCommandRequestProto readPutBlockRequest(ByteBuffer b) - throws IOException { - RatisHelper.debug(b, "readPutBlockRequest", LOG); - final ByteString byteString = ByteString.copyFrom(b); - - final ContainerCommandRequestProto request = - ContainerCommandRequestMessage.toProto(byteString, null); - - if (!request.hasPutBlock()) { - throw new StorageContainerException( - "Malformed PutBlock request. trace ID: " + request.getTraceID(), - Result.MALFORMED_REQUEST); - } - return request; - } - @Test public void testVolumeFullCase() throws Exception { File tempFile = File.createTempFile("test-kv-stream", ".tmp"); @@ -309,7 +325,7 @@ static ContainerCommandRequestProto closeBuffers( final ByteBuf buf = ref.retain(); final ContainerCommandRequestProto putBlockRequest; try { - putBlockRequest = readPutBlockRequest(buf); + putBlockRequest = KeyValueStreamDataChannel.readPutBlockRequest(buf); // write the remaining data writeFully(buf.nioBuffer(), writeMethod); } finally { @@ -401,4 +417,29 @@ static CompletableFuture completeExceptionally(Throwable t) { f.completeExceptionally(t); return f; } + + private static KeyValueStreamDataChannel newChannel(File tempFile) throws Exception { + HddsVolume mockVolume = mock(HddsVolume.class); + when(mockVolume.getStorageID()).thenReturn("storageId"); + when(mockVolume.getCurrentUsage()).thenReturn(new SpaceUsageSource.Fixed(1000L, 1000L, 0L)); + ContainerData mockContainerData = mock(ContainerData.class); + when(mockContainerData.getContainerID()).thenReturn(123L); + when(mockContainerData.getVolume()).thenReturn(mockVolume); + ContainerMetrics mockMetrics = mock(ContainerMetrics.class); + return new KeyValueStreamDataChannel(tempFile, mockContainerData, mockMetrics); + } + + private static void write(KeyValueStreamDataChannel channel, byte[] data) + throws IOException { + write(channel, ByteBuffer.wrap(data)); + } + + private static void write(KeyValueStreamDataChannel channel, ByteBuffer buf) + throws IOException { + ReferenceCountedObject ref = + ReferenceCountedObject.wrap(buf, () -> { }, () -> { }); + ref.retain(); + channel.write(ref); + ref.release(); + } } From b7ee889e3fb40e2e23dafd8c77d25fc390f48c18 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Thu, 9 Jul 2026 11:21:51 +0800 Subject: [PATCH 2/8] Remove unused import --- .../container/keyvalue/impl/TestKeyValueStreamDataChannel.java | 1 - 1 file changed, 1 deletion(-) diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index 08263a8c8f15..fd52ac389a81 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -26,7 +26,6 @@ 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.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; From de84fb82e28017e7a7835ad1e936ba106c9cc849 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Thu, 9 Jul 2026 13:37:33 +0800 Subject: [PATCH 3/8] fix wrong assert --- .../container/keyvalue/impl/KeyValueStreamDataChannel.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java index 73027e65ee81..65405d2ab666 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java @@ -253,7 +253,7 @@ private void closeWithStreamPutBlock() throws IOException { final ContainerCommandRequestProto proto = closeBuffers(buffers, super::writeFileChannel); putBlockRequest.set(proto); - Preconditions.checkState(putBlockProcessor == null); + Preconditions.checkState(putBlockProcessor != null); putBlockProcessor.processPutBlock(proto); setLinked(); } From 20d94737d9fb8b1e49cd575ac4a3484a95419976 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 21 Jul 2026 11:54:51 +0800 Subject: [PATCH 4/8] address comments --- .../container/common/impl/HddsDispatcher.java | 8 ++-- .../interfaces/ContainerDispatcher.java | 6 ++- .../container/common/interfaces/Handler.java | 4 +- .../server/ratis/ContainerStateMachine.java | 41 ++++++++----------- .../container/keyvalue/KeyValueHandler.java | 6 ++- .../keyvalue/impl/ChunkManagerDispatcher.java | 8 +++- .../keyvalue/impl/FilePerBlockStrategy.java | 8 ++-- .../impl/KeyValueStreamDataChannel.java | 31 +++++++------- .../keyvalue/impl/StreamDataChannelBase.java | 12 ++++++ .../keyvalue/interfaces/ChunkManager.java | 5 ++- .../impl/TestKeyValueStreamDataChannel.java | 13 ++---- .../main/proto/DatanodeClientProtocol.proto | 2 + 12 files changed, 81 insertions(+), 63 deletions(-) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java index 249e58df94d4..46d31f8156dd 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java @@ -80,6 +80,7 @@ import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; import org.apache.ratis.util.UncheckedAutoCloseable; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -807,13 +808,14 @@ private boolean isAllowed(String action) { @Override public StateMachine.DataChannel getStreamDataChannel( - ContainerCommandRequestProto msg) - throws StorageContainerException { + ContainerCommandRequestProto msg, + CheckedConsumer putBlock) + throws StorageContainerException { long containerID = msg.getContainerID(); Container container = getContainer(containerID); if (container != null) { Handler handler = getHandler(getContainerType(container)); - return handler.getStreamDataChannel(container, msg); + return handler.getStreamDataChannel(container, msg, putBlock); } else { throw new StorageContainerException( "ContainerID " + containerID + " does not exist", diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java index 2aba8253cefe..af6822443d8b 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java @@ -17,6 +17,7 @@ package org.apache.hadoop.ozone.container.common.interfaces; +import java.io.IOException; import java.util.Map; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; @@ -26,6 +27,7 @@ import org.apache.hadoop.ozone.container.common.transport.server.ratis.DispatcherContext; import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.function.CheckedConsumer; /** * Dispatcher acts as the bridge between the transport layer and @@ -87,7 +89,9 @@ void validateContainerCommand( * When uploading using stream, get StreamDataChannel. */ default StateMachine.DataChannel getStreamDataChannel( - ContainerCommandRequestProto msg) throws StorageContainerException { + ContainerCommandRequestProto msg, + CheckedConsumer putBlock) + throws StorageContainerException { throw new UnsupportedOperationException( "getStreamDataChannel not supported."); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java index fec625f89a08..6fe253ffbc7d 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java @@ -46,6 +46,7 @@ import org.apache.hadoop.ozone.container.keyvalue.TarContainerPacker; import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.function.CheckedConsumer; /** * Dispatcher sends ContainerCommandRequests to Handler. Each Container Type @@ -93,7 +94,8 @@ public static Handler getHandlerForContainerType( } public abstract StateMachine.DataChannel getStreamDataChannel( - Container container, ContainerCommandRequestProto msg) + Container container, ContainerCommandRequestProto msg, + CheckedConsumer putBlock) throws StorageContainerException; /** 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 3578df557f61..7a1129eff090 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 @@ -110,6 +110,7 @@ import org.apache.ratis.util.JavaUtils; import org.apache.ratis.util.LifeCycle; import org.apache.ratis.util.TaskQueue; +import org.apache.ratis.util.function.CheckedConsumer; import org.apache.ratis.util.function.CheckedSupplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -730,7 +731,17 @@ private StateMachine.DataChannel getStreamDataChannel( requestProto.getTraceID()); } dispatchCommand(requestProto, context); // stream init - return dispatcher.getStreamDataChannel(requestProto); + final CheckedConsumer putBlock + = requestProto.getCmdType() == Type.StreamInitWithPutBlock ? this::streamPutBlock : null; + return dispatcher.getStreamDataChannel(requestProto, putBlock); + } + + void streamPutBlock(ContainerCommandRequestProto request) throws IOException { + final DispatcherContext context = DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_LINK) + .setStage(DispatcherContext.WriteChunkStage.COMBINED) + .setContainer2BCSIDMap(container2BCSIDMap) + .build(); + dispatchCommand(request, context); } @Override @@ -745,16 +756,7 @@ public CompletableFuture stream(RaftClientRequest request) { .setStage(DispatcherContext.WriteChunkStage.WRITE_DATA) .setContainer2BCSIDMap(container2BCSIDMap) .build(); - DataChannel channel = getStreamDataChannel(requestProto, context); - if (datastreamPutBlockEnabled && channel instanceof KeyValueStreamDataChannel) { - final KeyValueStreamDataChannel kvChannel = (KeyValueStreamDataChannel) channel; - kvChannel.setDatastreamPutBlockEnabled(true); - kvChannel.setPutBlockProcessor(req -> dispatchCommand(req, - DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_LINK) - .setStage(DispatcherContext.WriteChunkStage.COMBINED) - .setContainer2BCSIDMap(container2BCSIDMap) - .build())); - } + final DataChannel channel = getStreamDataChannel(requestProto, context); final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ? getChunkExecutor(requestProto.getWriteChunk()) : null; return new LocalStream(channel, chunkExecutor); @@ -786,20 +788,11 @@ public CompletableFuture link(DataStream stream, LogEntryProto entry) { final KeyValueStreamDataChannel kvStreamDataChannel = (KeyValueStreamDataChannel) dataChannel; - - if (datastreamPutBlockEnabled) { - // The PutBlock should be commited when the stream is closed thus - // we expect this stream has been marked as linked. - if (kvStreamDataChannel.isLinked()) { - return CompletableFuture.completedFuture(null); - } else { - return JavaUtils.completeExceptionally(new IllegalStateException( - "PutBlock was not committed on stream close: " + kvStreamDataChannel)); - } - } else { - kvStreamDataChannel.setLinked(); - return CompletableFuture.completedFuture(null); + if (!kvStreamDataChannel.isLinked()) { + return JavaUtils.completeExceptionally(new IllegalStateException( + "PutBlock was not committed on stream close: " + kvStreamDataChannel)); } + return CompletableFuture.completedFuture(null); } private ExecutorService getChunkExecutor(WriteChunkRequestProto req) { 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 1b3f399b46fa..941f0479d4b0 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 @@ -173,6 +173,7 @@ import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -272,7 +273,8 @@ public KeyValueHandler(ConfigurationSource config, @Override public StateMachine.DataChannel getStreamDataChannel( - Container container, ContainerCommandRequestProto msg) + Container container, ContainerCommandRequestProto msg, + CheckedConsumer putBlock) throws StorageContainerException { KeyValueContainer kvContainer = (KeyValueContainer) container; checkContainerOpen(kvContainer); @@ -282,7 +284,7 @@ public StateMachine.DataChannel getStreamDataChannel( BlockID.getFromProtobuf(msg.getWriteChunk().getBlockID()); return chunkManager.getStreamDataChannel(kvContainer, - blockID, metrics); + blockID, putBlock, metrics); } else { throw new StorageContainerException("Malformed request.", ContainerProtos.Result.IO_EXCEPTION); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java index a83306ff79be..4cb634229d99 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java @@ -27,6 +27,7 @@ import java.util.Map; import java.util.Objects; import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.ozone.common.ChunkBuffer; import org.apache.hadoop.ozone.common.ChunkBufferToByteString; @@ -40,6 +41,7 @@ import org.apache.hadoop.ozone.container.keyvalue.interfaces.BlockManager; import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager; import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -79,10 +81,12 @@ public String streamInit(Container container, BlockID blockID) @Override public StateMachine.DataChannel getStreamDataChannel( - Container container, BlockID blockID, ContainerMetrics metrics) + Container container, BlockID blockID, + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { return selectHandler(container) - .getStreamDataChannel(container, blockID, metrics); + .getStreamDataChannel(container, blockID, putBlock, metrics); } @Override diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java index 9a13507d6b37..01000883f17e 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java @@ -58,6 +58,7 @@ import org.apache.hadoop.ozone.container.keyvalue.interfaces.BlockManager; import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager; import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -111,12 +112,13 @@ public String streamInit(Container container, BlockID blockID) @Override public StateMachine.DataChannel getStreamDataChannel( - Container container, BlockID blockID, ContainerMetrics metrics) - throws StorageContainerException { + Container container, BlockID blockID, + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { checkLayoutVersion(container); final File chunkFile = getChunkFile(container, blockID); return new KeyValueStreamDataChannel(chunkFile, - container.getContainerData(), metrics); + container.getContainerData(), putBlock, metrics); } @Override diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java index 65405d2ab666..9e86fd1293ee 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java @@ -35,6 +35,7 @@ import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf; import org.apache.ratis.util.ReferenceCountedObject; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -49,22 +50,13 @@ public class KeyValueStreamDataChannel extends StreamDataChannelBase { private final AtomicReference putBlockRequest = new AtomicReference<>(); private final AtomicBoolean closed = new AtomicBoolean(); - private boolean datastreamPutBlockEnabled; - private StreamPutBlockProcessor putBlockProcessor; + private final CheckedConsumer putBlock; KeyValueStreamDataChannel(File file, ContainerData containerData, - ContainerMetrics metrics) - throws StorageContainerException { + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { super(file, containerData, metrics); - datastreamPutBlockEnabled = false; - } - - public void setDatastreamPutBlockEnabled(boolean enabled) { - this.datastreamPutBlockEnabled = enabled; - } - - public void setPutBlockProcessor(StreamPutBlockProcessor processor) { - this.putBlockProcessor = processor; + this.putBlock = putBlock; } @Override @@ -72,6 +64,11 @@ ContainerProtos.Type getType() { return ContainerProtos.Type.StreamWrite; } + @Override + boolean isPutBlockCommittedOnClose() { + return putBlock == null; + } + @Override public int write(ReferenceCountedObject referenceCounted) throws IOException { @@ -120,7 +117,7 @@ void assertOpen() throws IOException { public void close() throws IOException { if (closed.compareAndSet(false, true)) { try { - if (datastreamPutBlockEnabled) { + if (isPutBlockCommittedOnClose()) { // This path requires the client appends the PutBlock at the // end of the stream. closeWithStreamPutBlock(); @@ -253,9 +250,9 @@ private void closeWithStreamPutBlock() throws IOException { final ContainerCommandRequestProto proto = closeBuffers(buffers, super::writeFileChannel); putBlockRequest.set(proto); - Preconditions.checkState(putBlockProcessor != null); - putBlockProcessor.processPutBlock(proto); - setLinked(); + Preconditions.checkState(putBlock != null); + putBlock.accept(proto); + link(); } interface WriteMethod { diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java index 3a65c3c42f5c..2cad1b052b21 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java @@ -94,6 +94,8 @@ public final boolean isOpen() { return getChannel().isOpen(); } + abstract boolean isPutBlockCommittedOnClose(); + protected void assertSpaceAvailability(int requested) throws StorageContainerException { ContainerUtils.assertSpaceAvailability(containerData.getContainerID(), containerData.getVolume(), requested); } @@ -102,6 +104,16 @@ public void setLinked() { linked.set(true); } + public boolean link() { + if (isPutBlockCommittedOnClose()) { + // The PutBlock should be commited when the steam is closed + return linked.get(); + } else { + setLinked(); + return true; + } + } + public boolean isLinked() { return linked.get(); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java index 0fc88a87ae3c..2cfcdf9c8e25 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java @@ -32,6 +32,7 @@ import org.apache.hadoop.ozone.container.common.transport.server.ratis.DispatcherContext; import org.apache.hadoop.ozone.container.keyvalue.KeyValueContainer; import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.util.function.CheckedConsumer; /** * Chunk Manager allows read, write, delete and listing of chunks in a container. @@ -114,7 +115,9 @@ default String streamInit(Container container, BlockID blockID) } default StateMachine.DataChannel getStreamDataChannel( - Container container, BlockID blockID, ContainerMetrics metrics) + Container container, BlockID blockID, + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { return null; } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index fd52ac389a81..f7de6e171eee 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -103,13 +103,8 @@ public class TestKeyValueStreamDataChannel { public void testClosePutBlockBehavior(boolean datastreamPutBlockEnabled) throws Exception { File tempFile = File.createTempFile("test-kv-close-" + datastreamPutBlockEnabled, ".tmp"); tempFile.deleteOnExit(); - KeyValueStreamDataChannel channel = newChannel(tempFile); + KeyValueStreamDataChannel channel = newChannelWithPutBlockCommittedOnClose(tempFile); AtomicReference processed = new AtomicReference<>(); - if (datastreamPutBlockEnabled) { - channel.setDatastreamPutBlockEnabled(datastreamPutBlockEnabled); - channel.setPutBlockProcessor(processed::set); - } - final byte[] data = RandomUtils.secure().randomBytes(50); final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( PUT_BLOCK_PROTO, null).getContent().asReadOnlyByteBuffer(); @@ -185,7 +180,7 @@ public void testVolumeFullCase() throws Exception { when(mockContainerData.getContainerID()).thenReturn(123L); when(mockContainerData.getVolume()).thenReturn(mockVolume); ContainerMetrics mockMetrics = mock(ContainerMetrics.class); - KeyValueStreamDataChannel writeChannel = new KeyValueStreamDataChannel(tempFile, mockContainerData, mockMetrics); + KeyValueStreamDataChannel writeChannel = new KeyValueStreamDataChannel(tempFile, mockContainerData, null, mockMetrics); assertThrows(StorageContainerException.class, () -> writeChannel.assertSpaceAvailability(1)); final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( @@ -417,7 +412,7 @@ static CompletableFuture completeExceptionally(Throwable t) { return f; } - private static KeyValueStreamDataChannel newChannel(File tempFile) throws Exception { + private static KeyValueStreamDataChannel newChannelWithPutBlockCommittedOnClose(File tempFile) throws Exception { HddsVolume mockVolume = mock(HddsVolume.class); when(mockVolume.getStorageID()).thenReturn("storageId"); when(mockVolume.getCurrentUsage()).thenReturn(new SpaceUsageSource.Fixed(1000L, 1000L, 0L)); @@ -425,7 +420,7 @@ private static KeyValueStreamDataChannel newChannel(File tempFile) throws Except when(mockContainerData.getContainerID()).thenReturn(123L); when(mockContainerData.getVolume()).thenReturn(mockVolume); ContainerMetrics mockMetrics = mock(ContainerMetrics.class); - return new KeyValueStreamDataChannel(tempFile, mockContainerData, mockMetrics); + return new KeyValueStreamDataChannel(tempFile, mockContainerData, null, mockMetrics); } private static void write(KeyValueStreamDataChannel channel, byte[] data) diff --git a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto index 05c94624c990..9438faafc7d9 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -88,6 +88,8 @@ enum Type { GetContainerChecksumInfo = 23; // Allows us to read a block ReadBlock = 24; + // Initializes a stream for writing data with PutBlock commited on close. + StreamInitWithPutBlock = 25; } From 780214198243ce8d41afce54c509255662443b34 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 21 Jul 2026 12:14:45 +0800 Subject: [PATCH 5/8] fix unit test after the big batch of changes. --- .../keyvalue/impl/KeyValueStreamDataChannel.java | 7 +++---- .../keyvalue/impl/TestKeyValueStreamDataChannel.java | 11 ++++++++--- 2 files changed, 11 insertions(+), 7 deletions(-) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java index 9e86fd1293ee..5d4f2ad82ad2 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java @@ -66,7 +66,7 @@ ContainerProtos.Type getType() { @Override boolean isPutBlockCommittedOnClose() { - return putBlock == null; + return putBlock != null; } @Override @@ -103,8 +103,7 @@ static void writeFully(ByteBuffer b, WriteMethod writeMethod) } public ContainerCommandRequestProto getPutBlockRequest() { - return Objects.requireNonNull(putBlockRequest.get(), - () -> "putBlockRequest == null, " + this); + return putBlockRequest.get(); } void assertOpen() throws IOException { @@ -252,7 +251,7 @@ private void closeWithStreamPutBlock() throws IOException { putBlockRequest.set(proto); Preconditions.checkState(putBlock != null); putBlock.accept(proto); - link(); + setLinked(); } interface WriteMethod { diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index f7de6e171eee..60e6267751ae 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -71,6 +71,7 @@ import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf; import org.apache.ratis.thirdparty.io.netty.buffer.Unpooled; import org.apache.ratis.util.ReferenceCountedObject; +import org.apache.ratis.util.function.CheckedConsumer; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; @@ -103,8 +104,8 @@ public class TestKeyValueStreamDataChannel { public void testClosePutBlockBehavior(boolean datastreamPutBlockEnabled) throws Exception { File tempFile = File.createTempFile("test-kv-close-" + datastreamPutBlockEnabled, ".tmp"); tempFile.deleteOnExit(); - KeyValueStreamDataChannel channel = newChannelWithPutBlockCommittedOnClose(tempFile); AtomicReference processed = new AtomicReference<>(); + KeyValueStreamDataChannel channel = newChannel(tempFile, datastreamPutBlockEnabled, processed); final byte[] data = RandomUtils.secure().randomBytes(50); final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( PUT_BLOCK_PROTO, null).getContent().asReadOnlyByteBuffer(); @@ -412,7 +413,9 @@ static CompletableFuture completeExceptionally(Throwable t) { return f; } - private static KeyValueStreamDataChannel newChannelWithPutBlockCommittedOnClose(File tempFile) throws Exception { + private static KeyValueStreamDataChannel newChannel( + File tempFile, boolean datastreamPutBlockEnabled, + AtomicReference processed) throws Exception { HddsVolume mockVolume = mock(HddsVolume.class); when(mockVolume.getStorageID()).thenReturn("storageId"); when(mockVolume.getCurrentUsage()).thenReturn(new SpaceUsageSource.Fixed(1000L, 1000L, 0L)); @@ -420,7 +423,9 @@ private static KeyValueStreamDataChannel newChannelWithPutBlockCommittedOnClose( when(mockContainerData.getContainerID()).thenReturn(123L); when(mockContainerData.getVolume()).thenReturn(mockVolume); ContainerMetrics mockMetrics = mock(ContainerMetrics.class); - return new KeyValueStreamDataChannel(tempFile, mockContainerData, null, mockMetrics); + CheckedConsumer putBlock = + datastreamPutBlockEnabled ? processed::set : null; + return new KeyValueStreamDataChannel(tempFile, mockContainerData, putBlock, mockMetrics); } private static void write(KeyValueStreamDataChannel channel, byte[] data) From 07f9777d6e8414853366a62073eac2bd84b8bc57 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 21 Jul 2026 12:33:35 +0800 Subject: [PATCH 6/8] Remove datanode side config --- .../statemachine/DatanodeConfiguration.java | 19 ------------------- .../server/ratis/ContainerStateMachine.java | 4 ---- .../TestDatanodeConfiguration.java | 17 ----------------- .../impl/TestKeyValueStreamDataChannel.java | 12 ++++++------ 4 files changed, 6 insertions(+), 46 deletions(-) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java index 276cbbd561d8..506dd79c37dc 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java @@ -535,17 +535,6 @@ public class DatanodeConfiguration extends ReconfigurableConfig { private boolean waitOnAllFollowers = WAIT_ON_ALL_FOLLOWERS_DEFAULT; - public static final String HDDS_DATANODE_DATASTREAM_PUTBLOCK_ENABLED = - CONFIG_PREFIX + ".datastream.putblock.enabled"; - - @Config(key = "hdds.datanode.datastream.putblock.enabled", - defaultValue = "false", - type = ConfigType.BOOLEAN, - tags = { DATANODE }, - description = "When enabled, PutBlock is committed when a Ratis data stream " + - "closes instead of via the Raft WriteAsync path.") - private boolean datastreamPutBlockEnabled = false; - @Config(key = "hdds.datanode.container.schema.v3.enabled", defaultValue = "true", type = ConfigType.BOOLEAN, @@ -1321,12 +1310,4 @@ public int getGrpcSoBacklog() { public void setGrpcSoBacklog(int grpcSoBacklog) { this.grpcSoBacklog = grpcSoBacklog; } - - public boolean isDatastreamPutBlockEnabled() { - return datastreamPutBlockEnabled; - } - - public void setDatastreamPutBlockEnabled(boolean datastreamPutBlockEnabled) { - this.datastreamPutBlockEnabled = datastreamPutBlockEnabled; - } } 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 7a1129eff090..8fd79b397f4c 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 @@ -161,7 +161,6 @@ public class ContainerStateMachine extends BaseStateMachine { private final Semaphore applyTransactionSemaphore; private final boolean waitOnBothFollowers; - private final boolean datastreamPutBlockEnabled; private final HddsDatanodeService datanodeService; private static Semaphore semaphore = new Semaphore(1); private final AtomicBoolean peersValidated; @@ -304,9 +303,6 @@ public ContainerStateMachine(HddsDatanodeService hddsDatanodeService, RaftGroupI this.waitOnBothFollowers = conf.getObject( DatanodeConfiguration.class).waitOnAllFollowers(); - this.datastreamPutBlockEnabled = conf.getObject( - DatanodeConfiguration.class).isDatastreamPutBlockEnabled(); - this.writeChunkWaitMaxNs = conf.getTimeDuration(ScmConfigKeys.HDDS_CONTAINER_RATIS_STATEMACHINE_WRITE_WAIT_INTERVAL, ScmConfigKeys.HDDS_CONTAINER_RATIS_STATEMACHINE_WRITE_WAIT_INTERVAL_NS_DEFAULT, TimeUnit.NANOSECONDS); } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java index 288a99095670..59e90ba30b38 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestDatanodeConfiguration.java @@ -359,21 +359,4 @@ void testGrpcSoBacklogSetter() { assertEquals(512, subject.getGrpcSoBacklog()); } - - @Test - void testDatastreamPutBlockEnabledDefault() { - DatanodeConfiguration subject = new OzoneConfiguration() - .getObject(DatanodeConfiguration.class); - assertThat(subject.isDatastreamPutBlockEnabled()).isFalse(); - } - - @Test - void testDatastreamPutBlockConfigParsing() { - OzoneConfiguration conf = new OzoneConfiguration(); - conf.setBoolean(DatanodeConfiguration.HDDS_DATANODE_DATASTREAM_PUTBLOCK_ENABLED, true); - - DatanodeConfiguration subject = conf.getObject(DatanodeConfiguration.class); - - assertThat(subject.isDatastreamPutBlockEnabled()).isTrue(); - } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index 60e6267751ae..978b46214403 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -101,11 +101,11 @@ public class TestKeyValueStreamDataChannel { @ParameterizedTest @ValueSource(booleans = {false, true}) - public void testClosePutBlockBehavior(boolean datastreamPutBlockEnabled) throws Exception { - File tempFile = File.createTempFile("test-kv-close-" + datastreamPutBlockEnabled, ".tmp"); + public void testClosePutBlockBehavior(boolean commitPutBlockOnClose) throws Exception { + File tempFile = File.createTempFile("test-kv-close-" + commitPutBlockOnClose, ".tmp"); tempFile.deleteOnExit(); AtomicReference processed = new AtomicReference<>(); - KeyValueStreamDataChannel channel = newChannel(tempFile, datastreamPutBlockEnabled, processed); + KeyValueStreamDataChannel channel = newChannel(tempFile, commitPutBlockOnClose, processed); final byte[] data = RandomUtils.secure().randomBytes(50); final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( PUT_BLOCK_PROTO, null).getContent().asReadOnlyByteBuffer(); @@ -119,7 +119,7 @@ public void testClosePutBlockBehavior(boolean datastreamPutBlockEnabled) throws assertThat(processed.get()).isNull(); channel.close(); - if (datastreamPutBlockEnabled) { + if (commitPutBlockOnClose) { assertEquals(PUT_BLOCK_PROTO, processed.get()); assertEquals(PUT_BLOCK_PROTO, channel.getPutBlockRequest()); assertTrue(channel.isLinked()); @@ -414,7 +414,7 @@ static CompletableFuture completeExceptionally(Throwable t) { } private static KeyValueStreamDataChannel newChannel( - File tempFile, boolean datastreamPutBlockEnabled, + File tempFile, boolean commitPutBlockOnClose, AtomicReference processed) throws Exception { HddsVolume mockVolume = mock(HddsVolume.class); when(mockVolume.getStorageID()).thenReturn("storageId"); @@ -424,7 +424,7 @@ private static KeyValueStreamDataChannel newChannel( when(mockContainerData.getVolume()).thenReturn(mockVolume); ContainerMetrics mockMetrics = mock(ContainerMetrics.class); CheckedConsumer putBlock = - datastreamPutBlockEnabled ? processed::set : null; + commitPutBlockOnClose ? processed::set : null; return new KeyValueStreamDataChannel(tempFile, mockContainerData, putBlock, mockMetrics); } From dd327d66df4bca259d119292e1f2a15138fb070b Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 21 Jul 2026 12:38:48 +0800 Subject: [PATCH 7/8] fix --- .../common/transport/server/ratis/ContainerStateMachine.java | 2 +- .../container/keyvalue/impl/KeyValueStreamDataChannel.java | 1 - .../container/keyvalue/impl/TestKeyValueStreamDataChannel.java | 3 ++- 3 files changed, 3 insertions(+), 3 deletions(-) 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 8fd79b397f4c..ef620986ea22 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 @@ -737,7 +737,7 @@ void streamPutBlock(ContainerCommandRequestProto request) throws IOException { .setStage(DispatcherContext.WriteChunkStage.COMBINED) .setContainer2BCSIDMap(container2BCSIDMap) .build(); - dispatchCommand(request, context); + dispatchCommand(request, context); } @Override diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java index 5d4f2ad82ad2..5a735812d796 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java @@ -21,7 +21,6 @@ import java.io.File; import java.io.IOException; import java.nio.ByteBuffer; -import java.util.Objects; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index 978b46214403..638c5707725e 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -181,7 +181,8 @@ public void testVolumeFullCase() throws Exception { when(mockContainerData.getContainerID()).thenReturn(123L); when(mockContainerData.getVolume()).thenReturn(mockVolume); ContainerMetrics mockMetrics = mock(ContainerMetrics.class); - KeyValueStreamDataChannel writeChannel = new KeyValueStreamDataChannel(tempFile, mockContainerData, null, mockMetrics); + KeyValueStreamDataChannel writeChannel = + new KeyValueStreamDataChannel(tempFile, mockContainerData, null, mockMetrics); assertThrows(StorageContainerException.class, () -> writeChannel.assertSpaceAvailability(1)); final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( From 988480dc847dacb8e8904b966482317c77abb3b3 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Wed, 22 Jul 2026 10:19:26 +0800 Subject: [PATCH 8/8] update --- .../server/ratis/ContainerStateMachine.java | 2 +- .../impl/StreamPutBlockProcessor.java | 29 ------------------- 2 files changed, 1 insertion(+), 30 deletions(-) delete mode 100644 hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java 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 ef620986ea22..9a3e988e27ad 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 @@ -784,7 +784,7 @@ public CompletableFuture link(DataStream stream, LogEntryProto entry) { final KeyValueStreamDataChannel kvStreamDataChannel = (KeyValueStreamDataChannel) dataChannel; - if (!kvStreamDataChannel.isLinked()) { + if (!kvStreamDataChannel.link()) { return JavaUtils.completeExceptionally(new IllegalStateException( "PutBlock was not committed on stream close: " + kvStreamDataChannel)); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java deleted file mode 100644 index fccf661bceed..000000000000 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamPutBlockProcessor.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * 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.keyvalue.impl; - -import java.io.IOException; -import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; - -/** - * Processes a PutBlock request received inline from the streaming write path. - */ -@FunctionalInterface -public interface StreamPutBlockProcessor { - void processPutBlock(ContainerCommandRequestProto request) throws IOException; -}