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..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,13 +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; + private boolean datastreamPutBlockWithoutRaftEnabled = false; @Config(key = "ozone.client.key.write.concurrency", defaultValue = "1", @@ -708,12 +711,12 @@ public void setStreamReadTimeout(Duration streamReadTimeout) { this.streamReadTimeout = streamReadTimeout; } - public boolean isDatastreamPutBlockOnCloseEnabled() { - return datastreamPutBlockOnCloseEnabled; + public boolean isDatastreamPutBlockWithoutRaftEnabled() { + return datastreamPutBlockWithoutRaftEnabled; } - public void setDatastreamPutBlockOnCloseEnabled(boolean datastreamPutBlockOnCloseEnabled) { - this.datastreamPutBlockOnCloseEnabled = datastreamPutBlockOnCloseEnabled; + 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 ae07ee09ca5c..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,9 +120,14 @@ public class BlockDataStreamOutput implements ByteBufferStreamOutput { // be released from the buffer pool. private final StreamCommitWatcher commitWatcher; + // 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; @@ -211,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 = @@ -322,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) { @@ -419,9 +427,9 @@ 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(); + waitPutBlockCommandFuturesComplete(); } final BlockData blockData = containerBlockData.build(); if (close) { @@ -444,11 +452,14 @@ 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.isDatastreamPutBlockWithoutRaftEnabled()) { + executePutBlockCommand(blockData, byteBufferList); + return; } try { XceiverClientReply asyncReply = @@ -498,6 +509,57 @@ 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); + putBlockCommandFutures.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); + } + // 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)); + } + + 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)); + } + public static CompletableFuture executePutBlockClose( ContainerCommandRequestProto putBlockRequest, int max, DataStreamOutput out) { @@ -559,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); @@ -596,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 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/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..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 @@ -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; @@ -78,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()); } } @@ -130,6 +131,58 @@ void writeFlushBoundaryTriggersPutBlock() throws Exception { assertEquals(2, pipeline.getReceivedPutBlocks().size()); } + @Test + void midStreamPutBlockUsesStreamCommandWhenEnabled() throws Exception { + MockDatanodePipeline pipeline = new MockDatanodePipeline(); + OzoneClientConfig config = createConfig(); + config.setDatastreamPutBlockWithoutRaftEnabled(true); + 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 add another PutBlock command: it is appended to the stream instead + assertEquals(1, pipeline.getReceivedCommands().size()); + assertEquals(0, pipeline.getReceivedPutBlocks().size()); + assertArrayEquals(data, pipeline.getAllReceivedData()); + } + + @Test + void midStreamPutBlockCommandReleasesBuffers() throws Exception { + MockDatanodePipeline pipeline = new MockDatanodePipeline(); + OzoneClientConfig config = createConfig(); + config.setDatastreamPutBlockWithoutRaftEnabled(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.setDatastreamPutBlockWithoutRaftEnabled(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..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 @@ -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; @@ -106,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; @@ -740,6 +742,33 @@ 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 { + // unsafeWrap is okay since the command buffer can be modified after this method has returned. + final ContainerCommandRequestProto request = ContainerCommandRequestMessage.toProto( + UnsafeByteOperations.unsafeWrap(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 +784,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..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 @@ -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,18 @@ public CompletableFuture cleanUp() { executor); } + @Override + 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); + } + }, 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..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 @@ -170,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; } @@ -210,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); } } @@ -267,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 = @@ -287,22 +287,30 @@ 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); + 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 { + assertThat(entry.getBlockID().getBlockCommitSequenceId()).isPositive(); + } validateData(client, keyName, data); } } @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 = @@ -329,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(); @@ -352,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);