From b8f8fc017e1e5d90c69c40a925b8dd3adab1a321 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 29 Sep 2026 12:35:14 +0800 Subject: [PATCH 1/6] HDDS-16008.PutBlocks from Flushes also go without Raft --- .../hadoop/hdds/scm/OzoneClientConfig.java | 18 +++++ .../scm/storage/BlockDataStreamOutput.java | 73 ++++++++++++++++++- .../hdds/scm/storage/StreamCommitWatcher.java | 10 ++- .../scm/storage/MockDatanodePipeline.java | 45 +++++++++++- .../storage/TestBlockDataStreamOutput.java | 69 ++++++++++++++++++ .../server/ratis/ContainerStateMachine.java | 29 +++++++- .../server/ratis/DispatcherContext.java | 4 +- .../transport/server/ratis/LocalStream.java | 21 +++++- .../common/impl/TestHddsDispatcher.java | 3 +- .../client/rpc/TestBlockDataStreamOutput.java | 32 ++++++++ 10 files changed, 294 insertions(+), 10 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java index 493d11f37bef..7a743f7377c8 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java @@ -325,6 +325,16 @@ public class OzoneClientConfig { tags = ConfigTag.CLIENT) private boolean datastreamPutBlockOnCloseEnabled = false; + @Config(key = "ozone.client.datastream.putblock.command.enabled", + defaultValue = "false", + type = ConfigType.BOOLEAN, + description = "When enabled, the PutBlock triggered by a flush in the middle of a Ratis data stream " + + "is sent as a data stream command instead of a separate WriteAsync PutBlock, " + + "so that it does not go through the Raft log. " + + "All the datanodes must support data stream commands.", + tags = ConfigTag.CLIENT) + private boolean datastreamPutBlockCommandEnabled = false; + @Config(key = "ozone.client.key.write.concurrency", defaultValue = "1", description = "Maximum concurrent writes allowed on each key. " + @@ -716,6 +726,14 @@ public void setDatastreamPutBlockOnCloseEnabled(boolean datastreamPutBlockOnClos this.datastreamPutBlockOnCloseEnabled = datastreamPutBlockOnCloseEnabled; } + public boolean isDatastreamPutBlockCommandEnabled() { + return datastreamPutBlockCommandEnabled; + } + + public void setDatastreamPutBlockCommandEnabled(boolean datastreamPutBlockCommandEnabled) { + this.datastreamPutBlockCommandEnabled = datastreamPutBlockCommandEnabled; + } + /** * Enum for indicating what mode to use when combining chunk and block * checksums to define an aggregate FileChecksum. This should be considered 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 ae07ee09ca5c..547e5077d755 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 @@ -29,6 +29,7 @@ import java.util.Queue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; +import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -120,8 +121,12 @@ public class BlockDataStreamOutput implements ByteBufferStreamOutput { // be released from the buffer pool. private final StreamCommitWatcher commitWatcher; - private Queue> - putBlockFutures = new LinkedList<>(); + private Queue> putBlockFutures = new LinkedList<>(); + + // Buffers acknowledged by a PutBlock committed through a data stream command. + // They are released by the caller thread since bufferList is not thread safe. + private final Queue> ackedBuffers + = new ConcurrentLinkedQueue<>(); private final List failedServers; private final Checksum checksum; @@ -379,6 +384,7 @@ public void writeOnRetry(long len) throws IOException { */ public void watchForCommit(boolean bufferFull) throws IOException { checkOpen(); + releaseAckedBuffers(); try { XceiverClientReply reply = bufferFull ? commitWatcher.watchOnFirstIndex() : @@ -449,6 +455,9 @@ public void executePutBlock(boolean close, // is no need to continue. return; } + } else if (config.isDatastreamPutBlockCommandEnabled()) { + executePutBlockCommand(blockData, byteBufferList); + return; } try { XceiverClientReply asyncReply = @@ -498,6 +507,66 @@ public void executePutBlock(boolean close, } } + /** + * Commit a PutBlock in the middle of the stream by sending it as a data stream command, + * so that it does not go through the Raft log. The command is ordered with the data written + * so far and the datanodes reply with the serialized {@link ContainerCommandResponseProto}. + * + * @param blockData the block metadata to commit + * @param byteBufferList the buffers covered by this PutBlock, to be released once it is acknowledged + */ + private void executePutBlockCommand(BlockData blockData, + List byteBufferList) throws IOException { + final ContainerCommandRequestProto putBlockRequest + = ContainerProtocolCalls.getPutBlockRequest( + xceiverClient.getPipeline(), blockData, false, tokenString); + final ByteBuffer command = ContainerCommandRequestMessage.toMessage( + putBlockRequest, null).getContent().asReadOnlyByteBuffer(); + Preconditions.checkState(command.remaining() <= PUT_BLOCK_REQUEST_LENGTH_MAX, + "PutBlock command length %s > max %s", command.remaining(), + PUT_BLOCK_REQUEST_LENGTH_MAX); + RatisHelper.debug(command, "putBlockCommand", LOG); + metrics.incrPendingContainerOpsMetrics(ContainerProtos.Type.PutBlock); + putBlockFutures.add(out.commandAsync(command) + .whenCompleteAsync((reply, e) -> { + metrics.decrPendingContainerOpsMetrics(ContainerProtos.Type.PutBlock); + try { + validatePutBlockCommandReply(reply, e); + } catch (IOException ioe) { + setIoException(ioe); + throw new CompletionException(ioe); + } + ackedBuffers.add(byteBufferList); + }, responseExecutor)); + } + + private void validatePutBlockCommandReply(DataStreamReply reply, Throwable e) + throws IOException { + if (e != null || reply == null || !reply.isSuccess()) { + throw new IOException("Failed to commit PutBlock for blockID " + blockID + + " through the data stream, reply=" + reply, e); + } + final ByteBuffer response = reply.nioBuffer(); + if (!response.hasRemaining()) { + throw new IOException("Datanodes in pipeline " + + xceiverClient.getPipeline().getId() + " did not handle the PutBlock" + + " command for blockID " + blockID + + "; they may be running a version which does not support it"); + } + validateResponse(ContainerCommandResponseProto.parseFrom(response)); + } + + /** + * Release the buffers of the PutBlock(s) committed through data stream commands. + * This is called by the caller thread since {@link #bufferList} is not thread safe. + */ + private void releaseAckedBuffers() { + for (List buffers = ackedBuffers.poll(); buffers != null; + buffers = ackedBuffers.poll()) { + commitWatcher.releaseBuffers(buffers); + } + } + public static CompletableFuture executePutBlockClose( ContainerCommandRequestProto putBlockRequest, int max, DataStreamOutput out) { diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java index d83ceae37d32..a8e4e478dff6 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java @@ -35,8 +35,16 @@ class StreamCommitWatcher extends AbstractCommitWatcher { @Override void releaseBuffers(long index) { + releaseBuffers(remove(index)); + } + + /** + * Release the given buffers, which have been acknowledged without a Raft log index, + * i.e. the PutBlock was committed by a data stream command. + */ + void releaseBuffers(List buffers) { long acked = 0; - for (StreamBuffer buffer : remove(index)) { + for (StreamBuffer buffer : buffers) { acked += buffer.position(); bufferList.remove(buffer); } diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/MockDatanodePipeline.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/MockDatanodePipeline.java index c1c11b0d1a83..c3cb97b7aa37 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/MockDatanodePipeline.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/MockDatanodePipeline.java @@ -77,8 +77,11 @@ public class MockDatanodePipeline { // Recorded state private final List receivedChunks = Collections.synchronizedList(new ArrayList<>()); private final List receivedPutBlocks = Collections.synchronizedList(new ArrayList<>()); + private final List receivedCommands = Collections.synchronizedList(new ArrayList<>()); private final AtomicInteger watchForCommitCount = new AtomicInteger(0); private volatile Type streamInitType; + // When true, stream commands are replied without a payload, as a datanode which does not handle them does. + private volatile boolean commandUnsupported; // Commit tracking private final AtomicLong nextLogIndex = new AtomicLong(1); @@ -202,6 +205,26 @@ public MockDatanodePipeline(BlockID blockID) throws IOException { return CompletableFuture.completedFuture(dataStreamReply(data.length)); }).when(mockDataStreamOutput).writeAsync(any(ByteBuffer.class), any(Iterable.class)); + // Setup commandAsync (mid-stream PutBlock) behavior + doAnswer(invocation -> { + ByteBuffer src = invocation.getArgument(0); + final ContainerCommandRequestProto request = toProto(src); + receivedCommands.add(request); + if (commandUnsupported) { + return CompletableFuture.completedFuture(dataStreamReply(0)); + } + int count = putBlockCount.incrementAndGet(); + if (count > putBlockFailAfter && putBlockFailure != null) { + CompletableFuture failed = new CompletableFuture<>(); + failed.completeExceptionally(putBlockFailure.get()); + return failed; + } + final DataStreamReply reply = dataStreamReply(0); + when(reply.nioBuffer()).thenReturn( + buildPutBlockResponse(blockID).toByteString().asReadOnlyByteBuffer()); + return CompletableFuture.completedFuture(reply); + }).when(mockDataStreamOutput).commandAsync(any(ByteBuffer.class)); + // Mock XceiverClientFactory this.clientFactory = mock(XceiverClientFactory.class); doReturn(xceiverClient).when(clientFactory).acquireClient(any(Pipeline.class), anyBoolean()); @@ -234,6 +257,11 @@ public List getReceivedPutBlocks() { return receivedPutBlocks; } + /** @return the requests received as data stream commands. */ + public List getReceivedCommands() { + return receivedCommands; + } + public int getWatchForCommitCount() { return watchForCommitCount.get(); } @@ -274,17 +302,26 @@ public MockDatanodePipeline failWatchAfter(int n, Supplier err) { return this; } + /** Reply to stream commands without a payload, as a datanode which does not handle them does. */ + public MockDatanodePipeline withUnsupportedCommand() { + this.commandUnsupported = true; + return this; + } + // --- Helpers --- private void captureStreamInitType(ByteBuffer buffer) { + streamInitType = toProto(buffer).getCmdType(); + } + + private static ContainerCommandRequestProto toProto(ByteBuffer buffer) { ByteBuffer dup = buffer.duplicate(); byte[] bytes = new byte[dup.remaining()]; dup.get(bytes); try { - streamInitType = ContainerCommandRequestMessage.toProto( - ByteString.copyFrom(bytes), null).getCmdType(); + return ContainerCommandRequestMessage.toProto(ByteString.copyFrom(bytes), null); } catch (Exception e) { - throw new IllegalStateException("Failed to decode stream init request", e); + throw new IllegalStateException("Failed to decode container command request", e); } } @@ -308,6 +345,8 @@ private static DataStreamReply dataStreamReply(long bytesWritten) { when(reply.getBytesWritten()).thenReturn(bytesWritten); when(reply.getDataLength()).thenReturn(bytesWritten); when(reply.getCommitInfos()).thenReturn(Collections.emptyList()); + // Ratis replies with an empty buffer unless the state machine returns a command reply. + when(reply.nioBuffer()).thenReturn(ByteBuffer.allocate(0).asReadOnlyBuffer()); return reply; } } diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java index 35aaac2a5af0..aa8eaffcc9bb 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java @@ -29,6 +29,7 @@ import java.util.List; import org.apache.commons.lang3.RandomUtils; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type; import org.apache.hadoop.hdds.scm.OzoneClientConfig; import org.junit.jupiter.api.Test; @@ -130,6 +131,74 @@ void writeFlushBoundaryTriggersPutBlock() throws Exception { assertEquals(2, pipeline.getReceivedPutBlocks().size()); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void midStreamPutBlockUsesStreamCommandWhenEnabled(boolean putBlockOnCloseEnabled) throws Exception { + MockDatanodePipeline pipeline = new MockDatanodePipeline(); + OzoneClientConfig config = createConfig(); + config.setDatastreamPutBlockCommandEnabled(true); + config.setDatastreamPutBlockOnCloseEnabled(putBlockOnCloseEnabled); + byte[] data = randomBytes(400); + try (BlockDataStreamOutput stream = createStream(pipeline, config)) { + stream.write(ByteBuffer.wrap(data), 0, data.length); + // 4 chunks of 100B each → hits flush boundary → 1 PutBlock, sent as a stream command + assertEquals(4, pipeline.getReceivedChunks().size()); + assertEquals(0, pipeline.getReceivedPutBlocks().size(), "The mid-stream PutBlock must not go through Raft"); + assertEquals(1, pipeline.getReceivedCommands().size()); + + ContainerCommandRequestProto command = pipeline.getReceivedCommands().get(0); + assertEquals(Type.PutBlock, command.getCmdType()); + assertThat(command.getPutBlock().getEof()).isFalse(); + assertEquals(4, command.getPutBlock().getBlockData().getChunksCount()); + } + // Close does not send another command; it commits PutBlock the way the other config selects. + assertEquals(1, pipeline.getReceivedCommands().size()); + assertEquals(putBlockOnCloseEnabled ? 0 : 1, pipeline.getReceivedPutBlocks().size()); + assertArrayEquals(data, pipeline.getAllReceivedData()); + } + + @Test + void midStreamPutBlockUsesRaftWhenCommandDisabled() throws Exception { + MockDatanodePipeline pipeline = new MockDatanodePipeline(); + OzoneClientConfig config = createConfig(); + config.setDatastreamPutBlockOnCloseEnabled(true); + byte[] data = randomBytes(400); + try (BlockDataStreamOutput stream = createStream(pipeline, config)) { + stream.write(ByteBuffer.wrap(data), 0, data.length); + assertEquals(1, pipeline.getReceivedPutBlocks().size(), "The mid-stream PutBlock should still go through Raft"); + assertEquals(0, pipeline.getReceivedCommands().size()); + } + assertArrayEquals(data, pipeline.getAllReceivedData()); + } + + @Test + void midStreamPutBlockCommandReleasesBuffers() throws Exception { + MockDatanodePipeline pipeline = new MockDatanodePipeline(); + OzoneClientConfig config = createConfig(); + config.setDatastreamPutBlockCommandEnabled(true); + BlockDataStreamOutput stream = createStream(pipeline, config); + byte[] data = randomBytes(400); + stream.write(ByteBuffer.wrap(data), 0, data.length); + stream.close(); + + assertEquals(400, stream.getTotalAckDataLength(), + "The buffers of a PutBlock committed by a stream command should be acknowledged"); + } + + @Test + void midStreamPutBlockFailsOnDatanodeWithoutCommandSupport() throws Exception { + MockDatanodePipeline pipeline = new MockDatanodePipeline().withUnsupportedCommand(); + OzoneClientConfig config = createConfig(); + config.setDatastreamPutBlockCommandEnabled(true); + BlockDataStreamOutput stream = createStream(pipeline, config); + byte[] data = randomBytes(400); + stream.write(ByteBuffer.wrap(data), 0, data.length); + + IOException e = assertThrows(IOException.class, stream::hsync); + assertThat(e).hasStackTraceContaining("did not handle the PutBlock command"); + assertThrows(IOException.class, stream::close); + } + @Test void writeAcrossStreamWindowTriggersBackPressure() throws Exception { MockDatanodePipeline pipeline = new MockDatanodePipeline(); 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 9a3e988e27ad..b4c4115a3e58 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 @@ -28,6 +28,7 @@ import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; +import java.nio.ByteBuffer; import java.nio.file.Files; import java.util.Arrays; import java.util.Collection; @@ -740,6 +741,32 @@ void streamPutBlock(ContainerCommandRequestProto request) throws IOException { dispatchCommand(request, context); } + /** + * Commit a PutBlock sent as a command in the middle of a data stream; + * see {@link StateMachine.DataStream#onCommand(ByteBuffer, long)}. + * The PutBlock is applied without a Raft log entry, so it has no log index to use as bcsId. + * + * @return the serialized {@link ContainerCommandResponseProto}, + * which is identical on all the peers of the pipeline. + */ + ByteBuffer streamCommand(ByteBuffer command) throws IOException { + final ContainerCommandRequestProto request = ContainerCommandRequestMessage.toProto( + ByteString.copyFrom(command), getGroupId()); + if (request.getCmdType() != Type.PutBlock) { + throw new StorageContainerException("Unexpected stream command " + request.getCmdType() + + ", expected " + Type.PutBlock, ContainerProtos.Result.MALFORMED_REQUEST); + } + final DispatcherContext context = DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_COMMAND) + .setStage(DispatcherContext.WriteChunkStage.COMBINED) + .setContainer2BCSIDMap(container2BCSIDMap) + .build(); + final ContainerCommandResponseProto response = dispatchCommand(request, context); + if (response.getResult() != ContainerProtos.Result.SUCCESS) { + throw new StorageContainerException(response.getMessage(), response.getResult()); + } + return response.toByteString().asReadOnlyByteBuffer(); + } + @Override public CompletableFuture stream(RaftClientRequest request) { return CompletableFuture.supplyAsync(() -> { @@ -755,7 +782,7 @@ public CompletableFuture stream(RaftClientRequest request) { final DataChannel channel = getStreamDataChannel(requestProto, context); final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ? getChunkExecutor(requestProto.getWriteChunk()) : null; - return new LocalStream(channel, chunkExecutor); + return new LocalStream(channel, chunkExecutor, this::streamCommand); } catch (IOException e) { throw new CompletionException("Failed to create data stream", e); } 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 e0f0ccbe193a..b5fd18aa4051 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 @@ -108,7 +108,8 @@ public enum Op { APPLY_TRANSACTION, STREAM_INIT, - STREAM_LINK; + STREAM_LINK, + STREAM_COMMAND; public boolean readFromTmpFile() { return this == READ_STATE_MACHINE_DATA; @@ -120,6 +121,7 @@ public boolean validateToken() { case WRITE_STATE_MACHINE_DATA: case READ_STATE_MACHINE_DATA: case STREAM_LINK: + case STREAM_COMMAND: return false; default: return true; diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java index f2b73e3e8fa1..95ea00201bea 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java @@ -17,19 +17,27 @@ package org.apache.hadoop.ozone.container.common.transport.server.ratis; +import java.io.IOException; +import java.nio.ByteBuffer; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.Executor; import org.apache.hadoop.ozone.container.keyvalue.impl.KeyValueStreamDataChannel; import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.util.JavaUtils; +import org.apache.ratis.util.function.CheckedFunction; class LocalStream implements StateMachine.DataStream { private final StateMachine.DataChannel dataChannel; private final Executor executor; + /** Handler for the commands sent in the middle of the stream. */ + private final CheckedFunction command; - LocalStream(StateMachine.DataChannel dataChannel, Executor executor) { + LocalStream(StateMachine.DataChannel dataChannel, Executor executor, + CheckedFunction command) { this.dataChannel = dataChannel; this.executor = executor; + this.command = command; } @Override @@ -48,6 +56,17 @@ public CompletableFuture cleanUp() { executor); } + @Override + public CompletableFuture onCommand(ByteBuffer buffer, long streamOffset) { + return CompletableFuture.supplyAsync(() -> { + try { + return command.apply(buffer); + } catch (IOException e) { + throw new CompletionException(e); + } + }, executor); + } + @Override public Executor getExecutor() { return executor; diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java index d60ca220ecec..f6137361ca97 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java @@ -907,7 +907,8 @@ public void verify(Token token, newContext(Op.WRITE_STATE_MACHINE_DATA, WriteChunkStage.WRITE_DATA), newContext(Op.READ_STATE_MACHINE_DATA), newContext(Op.APPLY_TRANSACTION), - newContext(Op.STREAM_LINK, WriteChunkStage.COMMIT_DATA) + newContext(Op.STREAM_LINK, WriteChunkStage.COMMIT_DATA), + newContext(Op.STREAM_COMMAND, WriteChunkStage.COMMIT_DATA) }; for (DispatcherContext context : notVerify) { LOG.info("notVerify {}", context); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java index c011c774c776..74a4a357273b 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java @@ -69,6 +69,7 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; +import org.junit.jupiter.params.provider.ValueSource; /** * Tests BlockDataStreamOutput class. @@ -298,6 +299,37 @@ public void testPutBlockAtBoundary(boolean flushDelay, boolean putBlockOnCloseEn } } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testPutBlockAtBoundaryByCommand(boolean putBlockOnCloseEnabled) throws Exception { + OzoneClientConfig config = newClientConfig(cluster.getConf(), false, putBlockOnCloseEnabled); + config.setDatastreamPutBlockCommandEnabled(true); + try (OzoneClient client = newClient(cluster.getConf(), config)) { + int dataLength = 500; + XceiverClientMetrics metrics = XceiverClientManager.getXceiverClientMetrics(); + long putBlockCount = metrics.getContainerOpCountMetrics(ContainerProtos.Type.PutBlock); + String keyName = getKeyName(); + OzoneDataStreamOutput key = createKey(client, keyName, 0); + byte[] data = ContainerTestHelper.getFixedLengthString(keyString, dataLength).getBytes(UTF_8); + key.write(ByteBuffer.wrap(data)); + BlockDataStreamOutputEntry entry = + ((KeyDataStreamOutput) key.getByteBufStreamOutput()).getStreamEntries().get(0); + key.close(); + // The PutBlock at the 400 byte flush boundary is sent as a data stream command; close adds another + // one only when it does not commit PutBlock through the stream. + int expectedPutBlocks = putBlockOnCloseEnabled ? 1 : 2; + assertEquals(metrics.getContainerOpCountMetrics(ContainerProtos.Type.PutBlock), + putBlockCount + expectedPutBlocks); + if (putBlockOnCloseEnabled) { + // No PutBlock went through Raft, so there is no log index to use as block commit sequence id. + assertEquals(0, entry.getBlockID().getBlockCommitSequenceId()); + } else { + assertThat(entry.getBlockID().getBlockCommitSequenceId()).isPositive(); + } + validateData(client, keyName, data); + } + } + @ParameterizedTest @MethodSource("clientParameters") public void testMinPacketSize(boolean flushDelay, boolean putBlockOnCloseEnabled) From 1d6d4521665f67ba82e1806ab32ceafec1ef7e58 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 29 Sep 2026 12:41:56 +0800 Subject: [PATCH 2/6] re-use the config --- .../hadoop/hdds/scm/OzoneClientConfig.java | 37 ++++------- .../scm/storage/BlockDataStreamOutput.java | 8 +-- .../hdds/scm/TestOzoneClientConfig.java | 10 +-- .../storage/TestBlockDataStreamOutput.java | 36 +++-------- .../client/rpc/TestBlockDataStreamOutput.java | 64 ++++++------------- .../rpc/TestContainerStateMachineStream.java | 10 +-- 6 files changed, 55 insertions(+), 110 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java index 7a743f7377c8..e88e62939077 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/OzoneClientConfig.java @@ -317,23 +317,16 @@ public class OzoneClientConfig { tags = ConfigTag.CLIENT) private boolean enablePutblockPiggybacking = false; - @Config(key = "ozone.client.datastream.putblock.on.close.enabled", + @Config(key = "ozone.client.datastream.putblock.without.raft.enabled", defaultValue = "false", type = ConfigType.BOOLEAN, - description = "When enabled, use StreamInitWithPutBlock so datanodes commit PutBlock " + - "when the Ratis data stream closes instead of via a separate WriteAsync PutBlock.", + description = "When enabled, the PutBlock of a Ratis data stream is committed through the data stream " + + "instead of a separate WriteAsync PutBlock, so that it does not go through the Raft log: " + + "StreamInitWithPutBlock is used so that datanodes commit PutBlock when the data stream closes, " + + "and the PutBlock triggered by a flush in the middle of the stream is sent as a data stream command. " + + "All the datanodes must support committing PutBlock through the data stream.", tags = ConfigTag.CLIENT) - private boolean datastreamPutBlockOnCloseEnabled = false; - - @Config(key = "ozone.client.datastream.putblock.command.enabled", - defaultValue = "false", - type = ConfigType.BOOLEAN, - description = "When enabled, the PutBlock triggered by a flush in the middle of a Ratis data stream " + - "is sent as a data stream command instead of a separate WriteAsync PutBlock, " + - "so that it does not go through the Raft log. " + - "All the datanodes must support data stream commands.", - tags = ConfigTag.CLIENT) - private boolean datastreamPutBlockCommandEnabled = false; + private boolean datastreamPutBlockWithoutRaftEnabled = false; @Config(key = "ozone.client.key.write.concurrency", defaultValue = "1", @@ -718,20 +711,12 @@ public void setStreamReadTimeout(Duration streamReadTimeout) { this.streamReadTimeout = streamReadTimeout; } - public boolean isDatastreamPutBlockOnCloseEnabled() { - return datastreamPutBlockOnCloseEnabled; - } - - public void setDatastreamPutBlockOnCloseEnabled(boolean datastreamPutBlockOnCloseEnabled) { - this.datastreamPutBlockOnCloseEnabled = datastreamPutBlockOnCloseEnabled; - } - - public boolean isDatastreamPutBlockCommandEnabled() { - return datastreamPutBlockCommandEnabled; + public boolean isDatastreamPutBlockWithoutRaftEnabled() { + return datastreamPutBlockWithoutRaftEnabled; } - public void setDatastreamPutBlockCommandEnabled(boolean datastreamPutBlockCommandEnabled) { - this.datastreamPutBlockCommandEnabled = datastreamPutBlockCommandEnabled; + public void setDatastreamPutBlockWithoutRaftEnabled(boolean datastreamPutBlockWithoutRaftEnabled) { + this.datastreamPutBlockWithoutRaftEnabled = datastreamPutBlockWithoutRaftEnabled; } /** 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 547e5077d755..6e3f8f2627e9 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 @@ -216,7 +216,7 @@ private DataStreamOutput setupStream(Pipeline pipeline) throws IOException { // TODO: The datanode UUID is not used meaningfully, consider deprecating // it or remove it completely if possible String id = pipeline.getFirstNode().getUuidString(); - ContainerProtos.Type streamInitType = config.isDatastreamPutBlockOnCloseEnabled() + ContainerProtos.Type streamInitType = config.isDatastreamPutBlockWithoutRaftEnabled() ? ContainerProtos.Type.StreamInitWithPutBlock : ContainerProtos.Type.StreamInit; ContainerProtos.ContainerCommandRequestProto.Builder builder = @@ -425,7 +425,7 @@ public void executePutBlock(boolean close, byteBufferList = null; } waitFuturesComplete(); - if (close && config.isDatastreamPutBlockOnCloseEnabled()) { + if (close && config.isDatastreamPutBlockWithoutRaftEnabled()) { // Wait for boundary PutBlock(s) before appending the stream-close PutBlock. waitPutBlockFuturesComplete(); } @@ -450,12 +450,12 @@ public void executePutBlock(boolean close, } } }); - if (config.isDatastreamPutBlockOnCloseEnabled()) { + if (config.isDatastreamPutBlockWithoutRaftEnabled()) { // PutBlock is supposed to be committed after the data stream close so there // is no need to continue. return; } - } else if (config.isDatastreamPutBlockCommandEnabled()) { + } else if (config.isDatastreamPutBlockWithoutRaftEnabled()) { executePutBlockCommand(blockData, byteBufferList); return; } diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/TestOzoneClientConfig.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/TestOzoneClientConfig.java index cb05d31a2d6e..c9ba90c87c85 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/TestOzoneClientConfig.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/TestOzoneClientConfig.java @@ -93,19 +93,19 @@ public void testStreamReadConfigParsing() { } @Test - void testDatastreamPutBlockOnCloseEnabledDefault() { + void testDatastreamPutBlockWithoutRaftEnabledDefault() { OzoneClientConfig subject = new OzoneConfiguration() .getObject(OzoneClientConfig.class); - assertFalse(subject.isDatastreamPutBlockOnCloseEnabled()); + assertFalse(subject.isDatastreamPutBlockWithoutRaftEnabled()); } @Test - void testDatastreamPutBlockOnCloseConfigParsing() { + void testDatastreamPutBlockWithoutRaftConfigParsing() { OzoneConfiguration conf = new OzoneConfiguration(); - conf.setBoolean("ozone.client.datastream.putblock.on.close.enabled", true); + conf.setBoolean("ozone.client.datastream.putblock.without.raft.enabled", true); OzoneClientConfig subject = conf.getObject(OzoneClientConfig.class); - assertTrue(subject.isDatastreamPutBlockOnCloseEnabled()); + assertTrue(subject.isDatastreamPutBlockWithoutRaftEnabled()); } } diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java index aa8eaffcc9bb..b9de99ad0168 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockDataStreamOutput.java @@ -79,12 +79,12 @@ private BlockDataStreamOutput createStream( @ParameterizedTest @ValueSource(booleans = {false, true}) - void streamInitTypeFollowsClientConfig(boolean putBlockOnCloseEnabled) throws Exception { + void streamInitTypeFollowsClientConfig(boolean putBlockWithoutRaft) throws Exception { MockDatanodePipeline pipeline = new MockDatanodePipeline(); OzoneClientConfig config = createConfig(); - config.setDatastreamPutBlockOnCloseEnabled(putBlockOnCloseEnabled); + config.setDatastreamPutBlockWithoutRaftEnabled(putBlockWithoutRaft); try (BlockDataStreamOutput stream = createStream(pipeline, config)) { - Type expected = putBlockOnCloseEnabled ? Type.StreamInitWithPutBlock : Type.StreamInit; + Type expected = putBlockWithoutRaft ? Type.StreamInitWithPutBlock : Type.StreamInit; assertEquals(expected, pipeline.getStreamInitType()); } } @@ -131,13 +131,11 @@ void writeFlushBoundaryTriggersPutBlock() throws Exception { assertEquals(2, pipeline.getReceivedPutBlocks().size()); } - @ParameterizedTest - @ValueSource(booleans = {false, true}) - void midStreamPutBlockUsesStreamCommandWhenEnabled(boolean putBlockOnCloseEnabled) throws Exception { + @Test + void midStreamPutBlockUsesStreamCommandWhenEnabled() throws Exception { MockDatanodePipeline pipeline = new MockDatanodePipeline(); OzoneClientConfig config = createConfig(); - config.setDatastreamPutBlockCommandEnabled(true); - config.setDatastreamPutBlockOnCloseEnabled(putBlockOnCloseEnabled); + config.setDatastreamPutBlockWithoutRaftEnabled(true); byte[] data = randomBytes(400); try (BlockDataStreamOutput stream = createStream(pipeline, config)) { stream.write(ByteBuffer.wrap(data), 0, data.length); @@ -151,23 +149,9 @@ void midStreamPutBlockUsesStreamCommandWhenEnabled(boolean putBlockOnCloseEnable assertThat(command.getPutBlock().getEof()).isFalse(); assertEquals(4, command.getPutBlock().getBlockData().getChunksCount()); } - // Close does not send another command; it commits PutBlock the way the other config selects. + // Close does not add another PutBlock command: it is appended to the stream instead assertEquals(1, pipeline.getReceivedCommands().size()); - assertEquals(putBlockOnCloseEnabled ? 0 : 1, pipeline.getReceivedPutBlocks().size()); - assertArrayEquals(data, pipeline.getAllReceivedData()); - } - - @Test - void midStreamPutBlockUsesRaftWhenCommandDisabled() throws Exception { - MockDatanodePipeline pipeline = new MockDatanodePipeline(); - OzoneClientConfig config = createConfig(); - config.setDatastreamPutBlockOnCloseEnabled(true); - byte[] data = randomBytes(400); - try (BlockDataStreamOutput stream = createStream(pipeline, config)) { - stream.write(ByteBuffer.wrap(data), 0, data.length); - assertEquals(1, pipeline.getReceivedPutBlocks().size(), "The mid-stream PutBlock should still go through Raft"); - assertEquals(0, pipeline.getReceivedCommands().size()); - } + assertEquals(0, pipeline.getReceivedPutBlocks().size()); assertArrayEquals(data, pipeline.getAllReceivedData()); } @@ -175,7 +159,7 @@ void midStreamPutBlockUsesRaftWhenCommandDisabled() throws Exception { void midStreamPutBlockCommandReleasesBuffers() throws Exception { MockDatanodePipeline pipeline = new MockDatanodePipeline(); OzoneClientConfig config = createConfig(); - config.setDatastreamPutBlockCommandEnabled(true); + config.setDatastreamPutBlockWithoutRaftEnabled(true); BlockDataStreamOutput stream = createStream(pipeline, config); byte[] data = randomBytes(400); stream.write(ByteBuffer.wrap(data), 0, data.length); @@ -189,7 +173,7 @@ void midStreamPutBlockCommandReleasesBuffers() throws Exception { void midStreamPutBlockFailsOnDatanodeWithoutCommandSupport() throws Exception { MockDatanodePipeline pipeline = new MockDatanodePipeline().withUnsupportedCommand(); OzoneClientConfig config = createConfig(); - config.setDatastreamPutBlockCommandEnabled(true); + config.setDatastreamPutBlockWithoutRaftEnabled(true); BlockDataStreamOutput stream = createStream(pipeline, config); byte[] data = randomBytes(400); stream.write(ByteBuffer.wrap(data), 0, data.length); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java index 74a4a357273b..6a431e477b98 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestBlockDataStreamOutput.java @@ -69,7 +69,6 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; -import org.junit.jupiter.params.provider.ValueSource; /** * Tests BlockDataStreamOutput class. @@ -171,17 +170,17 @@ private static Stream dataLengthParameters() { private static Stream streamWriteParameters() { return dataLengthParameters().flatMap(dataLength -> - Stream.of(true, false).map(putBlockOnCloseEnabled -> - Arguments.of(dataLength.get()[0], putBlockOnCloseEnabled))); + Stream.of(true, false).map(putBlockWithoutRaft -> + Arguments.of(dataLength.get()[0], putBlockWithoutRaft))); } static OzoneClientConfig newClientConfig(ConfigurationSource source, boolean flushDelay, - boolean putBlockOnCloseEnabled) { + boolean putBlockWithoutRaft) { OzoneClientConfig clientConfig = source.getObject(OzoneClientConfig.class); clientConfig.setChecksumType(ContainerProtos.ChecksumType.NONE); clientConfig.setStreamBufferFlushDelay(flushDelay); - clientConfig.setDatastreamPutBlockOnCloseEnabled(putBlockOnCloseEnabled); + clientConfig.setDatastreamPutBlockWithoutRaftEnabled(putBlockWithoutRaft); return clientConfig; } @@ -211,13 +210,13 @@ public void shutdown() { @ParameterizedTest @MethodSource("streamWriteParameters") @Flaky("HDDS-12027") - public void testStreamWrite(int dataLength, boolean putBlockOnCloseEnabled) throws Exception { - OzoneClientConfig config = newClientConfig(cluster.getConf(), false, putBlockOnCloseEnabled); + public void testStreamWrite(int dataLength, boolean putBlockWithoutRaft) throws Exception { + OzoneClientConfig config = newClientConfig(cluster.getConf(), false, putBlockWithoutRaft); try (OzoneClient client = newClient(cluster.getConf(), config)) { testWrite(client, dataLength); // Forced container close before stream close relies on async PutBlock recovery; // that path is not used when PutBlock is committed only on data stream close. - if (!putBlockOnCloseEnabled) { + if (!putBlockWithoutRaft) { testWriteWithFailure(client, dataLength); } } @@ -268,9 +267,9 @@ static void validateData(OzoneClient client, String keyName, byte[] data) throws @ParameterizedTest @MethodSource("clientParameters") - public void testPutBlockAtBoundary(boolean flushDelay, boolean putBlockOnCloseEnabled) + public void testPutBlockAtBoundary(boolean flushDelay, boolean putBlockWithoutRaft) throws Exception { - OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockOnCloseEnabled); + OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockWithoutRaft); try (OzoneClient client = newClient(cluster.getConf(), config)) { int dataLength = 500; XceiverClientMetrics metrics = @@ -288,39 +287,16 @@ public void testPutBlockAtBoundary(boolean flushDelay, boolean putBlockOnCloseEn key.write(ByteBuffer.wrap(data)); assertThat(metrics.getPendingContainerOpCountMetrics(ContainerProtos.Type.PutBlock)) .isLessThanOrEqualTo(pendingPutBlockCount + 1); + BlockDataStreamOutputEntry entry = + ((KeyDataStreamOutput) key.getByteBufStreamOutput()).getStreamEntries().get(0); key.close(); // Since data length is 500, first putBlock will be at 400 (flush boundary). - // Close commits via WriteAsync PutBlock only when putBlockOnClose is disabled. - int expectedPutBlocks = putBlockOnCloseEnabled ? 1 : 2; + // Close commits via WriteAsync PutBlock only when PutBlock does not go through the data stream. + int expectedPutBlocks = putBlockWithoutRaft ? 1 : 2; assertEquals( metrics.getContainerOpCountMetrics(ContainerProtos.Type.PutBlock), putBlockCount + expectedPutBlocks); - validateData(client, keyName, data); - } - } - - @ParameterizedTest - @ValueSource(booleans = {false, true}) - public void testPutBlockAtBoundaryByCommand(boolean putBlockOnCloseEnabled) throws Exception { - OzoneClientConfig config = newClientConfig(cluster.getConf(), false, putBlockOnCloseEnabled); - config.setDatastreamPutBlockCommandEnabled(true); - try (OzoneClient client = newClient(cluster.getConf(), config)) { - int dataLength = 500; - XceiverClientMetrics metrics = XceiverClientManager.getXceiverClientMetrics(); - long putBlockCount = metrics.getContainerOpCountMetrics(ContainerProtos.Type.PutBlock); - String keyName = getKeyName(); - OzoneDataStreamOutput key = createKey(client, keyName, 0); - byte[] data = ContainerTestHelper.getFixedLengthString(keyString, dataLength).getBytes(UTF_8); - key.write(ByteBuffer.wrap(data)); - BlockDataStreamOutputEntry entry = - ((KeyDataStreamOutput) key.getByteBufStreamOutput()).getStreamEntries().get(0); - key.close(); - // The PutBlock at the 400 byte flush boundary is sent as a data stream command; close adds another - // one only when it does not commit PutBlock through the stream. - int expectedPutBlocks = putBlockOnCloseEnabled ? 1 : 2; - assertEquals(metrics.getContainerOpCountMetrics(ContainerProtos.Type.PutBlock), - putBlockCount + expectedPutBlocks); - if (putBlockOnCloseEnabled) { + if (putBlockWithoutRaft) { // No PutBlock went through Raft, so there is no log index to use as block commit sequence id. assertEquals(0, entry.getBlockID().getBlockCommitSequenceId()); } else { @@ -332,9 +308,9 @@ public void testPutBlockAtBoundaryByCommand(boolean putBlockOnCloseEnabled) thro @ParameterizedTest @MethodSource("clientParameters") - public void testMinPacketSize(boolean flushDelay, boolean putBlockOnCloseEnabled) + public void testMinPacketSize(boolean flushDelay, boolean putBlockWithoutRaft) throws Exception { - OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockOnCloseEnabled); + OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockWithoutRaft); try (OzoneClient client = newClient(cluster.getConf(), config)) { String keyName = getKeyName(); XceiverClientMetrics metrics = @@ -361,9 +337,9 @@ public void testMinPacketSize(boolean flushDelay, boolean putBlockOnCloseEnabled @ParameterizedTest @MethodSource("clientParameters") - public void testTotalAckDataLength(boolean flushDelay, boolean putBlockOnCloseEnabled) + public void testTotalAckDataLength(boolean flushDelay, boolean putBlockWithoutRaft) throws Exception { - OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockOnCloseEnabled); + OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockWithoutRaft); try (OzoneClient client = newClient(cluster.getConf(), config)) { int dataLength = 400; String keyName = getKeyName(); @@ -384,9 +360,9 @@ public void testTotalAckDataLength(boolean flushDelay, boolean putBlockOnCloseEn @ParameterizedTest @MethodSource("clientParameters") - public void testDatanodeVersion(boolean flushDelay, boolean putBlockOnCloseEnabled) + public void testDatanodeVersion(boolean flushDelay, boolean putBlockWithoutRaft) throws Exception { - OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockOnCloseEnabled); + OzoneClientConfig config = newClientConfig(cluster.getConf(), flushDelay, putBlockWithoutRaft); try (OzoneClient client = newClient(cluster.getConf(), config)) { // Verify all DNs internally have versions set correctly List dns = cluster.getHddsDatanodes(); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestContainerStateMachineStream.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestContainerStateMachineStream.java index 5e10814f5080..c10dfe784fb3 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestContainerStateMachineStream.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestContainerStateMachineStream.java @@ -77,23 +77,23 @@ void shutdown() { private static Stream streamingParameters() { return Stream.of(-1, +1).flatMap(offset -> - Stream.of(false, true).map(putBlockOnCloseEnabled -> - Arguments.of(offset, putBlockOnCloseEnabled))); + Stream.of(false, true).map(putBlockWithoutRaft -> + Arguments.of(offset, putBlockWithoutRaft))); } @ParameterizedTest @MethodSource("streamingParameters") - void testContainerStateMachineForStreaming(int offset, boolean putBlockOnCloseEnabled) + void testContainerStateMachineForStreaming(int offset, boolean putBlockWithoutRaft) throws Exception { final int size = chunkSize + offset; OzoneConfiguration conf = new OzoneConfiguration(cluster().getConf()); OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class); - clientConfig.setDatastreamPutBlockOnCloseEnabled(putBlockOnCloseEnabled); + clientConfig.setDatastreamPutBlockWithoutRaftEnabled(putBlockWithoutRaft); conf.setFromObject(clientConfig); final List locationInfoList; try (OzoneClient streamingClient = OzoneClientFactory.getRpcClient(conf); - OzoneDataStreamOutput key = createStreamKey("key" + offset + "-" + putBlockOnCloseEnabled, + OzoneDataStreamOutput key = createStreamKey("key" + offset + "-" + putBlockWithoutRaft, ReplicationType.RATIS, size, streamingClient.getObjectStore(), volumeName, bucketName)) { byte[] data = ContainerTestHelper.generateData(size, true); From b7aac57a382ba4ef74ba1bebc8fcad577cca9011 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 29 Sep 2026 14:41:54 +0800 Subject: [PATCH 3/6] re-use the old path to clean up committed buffers. --- .../scm/storage/BlockDataStreamOutput.java | 22 +++---------------- .../hdds/scm/storage/StreamCommitWatcher.java | 10 +-------- 2 files changed, 4 insertions(+), 28 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 6e3f8f2627e9..3b2fc3a5715d 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 @@ -29,7 +29,6 @@ import java.util.Queue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; -import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -123,11 +122,6 @@ public class BlockDataStreamOutput implements ByteBufferStreamOutput { private Queue> putBlockFutures = new LinkedList<>(); - // Buffers acknowledged by a PutBlock committed through a data stream command. - // They are released by the caller thread since bufferList is not thread safe. - private final Queue> ackedBuffers - = new ConcurrentLinkedQueue<>(); - private final List failedServers; private final Checksum checksum; @@ -384,7 +378,6 @@ public void writeOnRetry(long len) throws IOException { */ public void watchForCommit(boolean bufferFull) throws IOException { checkOpen(); - releaseAckedBuffers(); try { XceiverClientReply reply = bufferFull ? commitWatcher.watchOnFirstIndex() : @@ -536,7 +529,9 @@ private void executePutBlockCommand(BlockData blockData, setIoException(ioe); throw new CompletionException(ioe); } - ackedBuffers.add(byteBufferList); + // The command has no log index; use 0 as for the standalone protocol + // so that the buffers are released by the next watchForCommit. + commitWatcher.updateCommitInfoMap(0, byteBufferList); }, responseExecutor)); } @@ -556,17 +551,6 @@ private void validatePutBlockCommandReply(DataStreamReply reply, Throwable e) validateResponse(ContainerCommandResponseProto.parseFrom(response)); } - /** - * Release the buffers of the PutBlock(s) committed through data stream commands. - * This is called by the caller thread since {@link #bufferList} is not thread safe. - */ - private void releaseAckedBuffers() { - for (List buffers = ackedBuffers.poll(); buffers != null; - buffers = ackedBuffers.poll()) { - commitWatcher.releaseBuffers(buffers); - } - } - public static CompletableFuture executePutBlockClose( ContainerCommandRequestProto putBlockRequest, int max, DataStreamOutput out) { diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java index a8e4e478dff6..d83ceae37d32 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamCommitWatcher.java @@ -35,16 +35,8 @@ class StreamCommitWatcher extends AbstractCommitWatcher { @Override void releaseBuffers(long index) { - releaseBuffers(remove(index)); - } - - /** - * Release the given buffers, which have been acknowledged without a Raft log index, - * i.e. the PutBlock was committed by a data stream command. - */ - void releaseBuffers(List buffers) { long acked = 0; - for (StreamBuffer buffer : buffers) { + for (StreamBuffer buffer : remove(index)) { acked += buffer.position(); bufferList.remove(buffer); } From 93b3f105ec3d0f3aab0865e98ba66b5a95a906ab Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 29 Sep 2026 16:08:56 +0800 Subject: [PATCH 4/6] seperate future queues --- .../scm/storage/BlockDataStreamOutput.java | 22 ++++++++++++++----- 1 file changed, 16 insertions(+), 6 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 3b2fc3a5715d..e8205859395a 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 @@ -120,7 +120,13 @@ public class BlockDataStreamOutput implements ByteBufferStreamOutput { // be released from the buffer pool. private final StreamCommitWatcher commitWatcher; - private Queue> putBlockFutures = new LinkedList<>(); + // Futures of the PutBlock(s) sent through the Ratis write API, i.e. committed by the Raft log. + private Queue> + putBlockFutures = new LinkedList<>(); + + // Futures of the PutBlock(s) sent as data stream commands, i.e. committed without the Raft log. + private final Queue> + putBlockCommandFutures = new LinkedList<>(); private final List failedServers; private final Checksum checksum; @@ -321,6 +327,9 @@ private void doFlushIfNeeded() throws IOException { if (!putBlockFutures.isEmpty()) { putBlockFutures.remove().get(); } + if (!putBlockCommandFutures.isEmpty()) { + putBlockCommandFutures.remove().get(); + } } catch (ExecutionException e) { handleExecutionException(e); } catch (InterruptedException ex) { @@ -420,7 +429,7 @@ public void executePutBlock(boolean close, waitFuturesComplete(); if (close && config.isDatastreamPutBlockWithoutRaftEnabled()) { // Wait for boundary PutBlock(s) before appending the stream-close PutBlock. - waitPutBlockFuturesComplete(); + waitPutBlockCommandFuturesComplete(); } final BlockData blockData = containerBlockData.build(); if (close) { @@ -520,7 +529,7 @@ private void executePutBlockCommand(BlockData blockData, PUT_BLOCK_REQUEST_LENGTH_MAX); RatisHelper.debug(command, "putBlockCommand", LOG); metrics.incrPendingContainerOpsMetrics(ContainerProtos.Type.PutBlock); - putBlockFutures.add(out.commandAsync(command) + putBlockCommandFutures.add(out.commandAsync(command) .whenCompleteAsync((reply, e) -> { metrics.decrPendingContainerOpsMetrics(ContainerProtos.Type.PutBlock); try { @@ -612,12 +621,12 @@ public void waitFuturesComplete() throws IOException { } } - private void waitPutBlockFuturesComplete() throws IOException { - if (putBlockFutures.isEmpty()) { + private void waitPutBlockCommandFuturesComplete() throws IOException { + if (putBlockCommandFutures.isEmpty()) { return; } try { - CompletableFuture.allOf(putBlockFutures.toArray(EMPTY_FUTURE_ARRAY)).get(); + CompletableFuture.allOf(putBlockCommandFutures.toArray(EMPTY_FUTURE_ARRAY)).get(); checkOpen(); } catch (Exception e) { LOG.warn("Failed to commit PutBlock before stream close: " + e); @@ -649,6 +658,7 @@ private void handleFlush(boolean close) executePutBlock(true, true); } CompletableFuture.allOf(putBlockFutures.toArray(EMPTY_FUTURE_ARRAY)).get(); + CompletableFuture.allOf(putBlockCommandFutures.toArray(EMPTY_FUTURE_ARRAY)).get(); watchForCommit(false); // just check again if the exception is hit while waiting for the // futures to ensure flush has indeed succeeded From 0eb631266918eade2090cd624c725538aa8251a2 Mon Sep 17 00:00:00 2001 From: amaliujia Date: Tue, 6 Oct 2026 10:42:13 +0800 Subject: [PATCH 5/6] update --- .../common/transport/server/ratis/ContainerStateMachine.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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 b4c4115a3e58..f77bfadefad9 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 @@ -107,6 +107,7 @@ import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.com.google.protobuf.InvalidProtocolBufferException; import org.apache.ratis.thirdparty.com.google.protobuf.TextFormat; +import org.apache.ratis.thirdparty.com.google.protobuf.UnsafeByteOperations; import org.apache.ratis.util.FileUtils; import org.apache.ratis.util.JavaUtils; import org.apache.ratis.util.LifeCycle; @@ -750,8 +751,9 @@ void streamPutBlock(ContainerCommandRequestProto request) throws IOException { * which is identical on all the peers of the pipeline. */ ByteBuffer streamCommand(ByteBuffer command) throws IOException { + // unsafeWrap is okay since the command buffer can be modified after this method has returned. final ContainerCommandRequestProto request = ContainerCommandRequestMessage.toProto( - ByteString.copyFrom(command), getGroupId()); + UnsafeByteOperations.unsafeWrap(command), getGroupId()); if (request.getCmdType() != Type.PutBlock) { throw new StorageContainerException("Unexpected stream command " + request.getCmdType() + ", expected " + Type.PutBlock, ContainerProtos.Result.MALFORMED_REQUEST); From 0a42d071b9c7e9af1a8ba7b9f20ed14439caaf83 Mon Sep 17 00:00:00 2001 From: Rui Wang Date: Tue, 6 Oct 2026 16:06:42 +0800 Subject: [PATCH 6/6] Update hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java Co-authored-by: Peter Lee --- .../container/common/transport/server/ratis/LocalStream.java | 1 + 1 file changed, 1 insertion(+) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java index 95ea00201bea..47189031a11f 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/LocalStream.java @@ -60,6 +60,7 @@ public CompletableFuture cleanUp() { public CompletableFuture onCommand(ByteBuffer buffer, long streamOffset) { return CompletableFuture.supplyAsync(() -> { try { + // TODO: drain datachannel.buffers return command.apply(buffer); } catch (IOException e) { throw new CompletionException(e);