Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -317,13 +317,16 @@ public class OzoneClientConfig {
tags = ConfigTag.CLIENT)
private boolean enablePutblockPiggybacking = false;

@Config(key = "ozone.client.datastream.putblock.on.close.enabled",

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should add a new conf but not changing an existing conf. Otherwise, it becomes an incompatible change.

@amaliujia amaliujia Oct 5, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This feature is not released at all (for example in 2.1.2, there is no such code yet), so there is no compatibility that we need to maintain.

Then it is better that we only use one config for simplification.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, you are right! Then, we should change it instead.

@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",
Expand Down Expand Up @@ -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;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<CompletableFuture<ContainerCommandResponseProto>>
putBlockFutures = new LinkedList<>();

// Futures of the PutBlock(s) sent as data stream commands, i.e. committed without the Raft log.
private final Queue<CompletableFuture<DataStreamReply>>
putBlockCommandFutures = new LinkedList<>();

private final List<DatanodeDetails> failedServers;
private final Checksum checksum;

Expand Down Expand Up @@ -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 =
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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) {
Expand All @@ -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 =
Expand Down Expand Up @@ -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<StreamBuffer> 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<DataStreamReply> executePutBlockClose(
ContainerCommandRequestProto putBlockRequest, int max,
DataStreamOutput out) {
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -77,8 +77,11 @@ public class MockDatanodePipeline {
// Recorded state
private final List<byte[]> receivedChunks = Collections.synchronizedList(new ArrayList<>());
private final List<ContainerCommandRequestProto> receivedPutBlocks = Collections.synchronizedList(new ArrayList<>());
private final List<ContainerCommandRequestProto> 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);
Expand Down Expand Up @@ -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<DataStreamReply> 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());
Expand Down Expand Up @@ -234,6 +257,11 @@ public List<ContainerCommandRequestProto> getReceivedPutBlocks() {
return receivedPutBlocks;
}

/** @return the requests received as data stream commands. */
public List<ContainerCommandRequestProto> getReceivedCommands() {
return receivedCommands;
}

public int getWatchForCommitCount() {
return watchForCommitCount.get();
}
Expand Down Expand Up @@ -274,17 +302,26 @@ public MockDatanodePipeline failWatchAfter(int n, Supplier<Throwable> 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);
}
}

Expand All @@ -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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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());
}
}
Expand Down Expand Up @@ -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();
Expand Down
Loading
Loading