Skip to content

HDDS-16008. PutBlocks from Flushes also go without Raft - #11356

Merged
amaliujia merged 6 commits into
apache:masterfrom
amaliujia:putblock_on_flush_without_raft
Oct 6, 2026
Merged

amaliujia merged 6 commits into
apache:masterfrom
amaliujia:putblock_on_flush_without_raft

Conversation

@amaliujia

@amaliujia amaliujia commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

Apply https://issues.apache.org/jira/browse/RATIS-2625 to Ozone PutBlock commit. Use the new added Ratis data stream command API to commit PutBlock during flushes (which typically happens before the stream closes).

By doing so, once the config is enabled, all the PutBlock commits will be raftless.

What is the link to the Apache JIRA

https://issues.apache.org/jira/browse/HDDS-16008

How was this patch tested?

Unit test and integration test

@amaliujia amaliujia changed the title HDDS-16008.PutBlocks from Flushes also go without Raft HDDS-16008. PutBlocks from Flushes also go without Raft Sep 29, 2026
@amaliujia
amaliujia requested a review from szetszwo September 30, 2026 03:44
@amaliujia

Copy link
Copy Markdown
Contributor Author

The core changes are minors. Majority LOC is about testing.

@szetszwo szetszwo left a comment

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.

@amaliujia , thanks for working on this! Please see the comments inlined.

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.

Comment on lines +541 to +543
// 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);

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.

Let's just release the buffer here without using commitWatcher?

@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.

done. though I need to add a lock to byteBufferList in this case. You can tell from the diff how much mess I need to introduce now.

I feel like re-use commitWatcher.updateCommitInfoMap is still fine and have less complexity.

@szetszwo szetszwo Oct 5, 2026 •

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.

Then, let's use commitWatcher for now. We may simplify the code later on. Sorry that I was aware of the complication.

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.

Sure revert back to use commitWatcher

Comment on lines +753 to +754
final ContainerCommandRequestProto request = ContainerCommandRequestMessage.toProto(
ByteString.copyFrom(command), getGroupId());

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.

Add a new ContainerCommandRequestMessage.toProto(ByteBuffer, RaftGroupId) method to avoid copying.

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.

done

@szetszwo szetszwo Oct 5, 2026 •

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 may unsafeWrap(..) in streamCommand.

  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());

@szetszwo

szetszwo commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

@amaliujia , thanks for the update! Please see the responses above.

@amaliujia
amaliujia force-pushed the putblock_on_flush_without_raft branch from 2df1f8b to 0eb6312 Compare October 6, 2026 02:42

@szetszwo szetszwo left a comment

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.

+1 the change looks good.

* @return the serialized {@link ContainerCommandResponseProto},
* which is identical on all the peers of the pipeline.
*/
ByteBuffer streamCommand(ByteBuffer command) throws IOException {

@peterxcli peterxcli Oct 6, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

nit: wish we can refactor container SM and local stream, so we dont need to pass function as lambda. one reference-only refactor is -- LocalStream becomes the one place that handles a PutBlock arriving on the stream.

public class ContainerStateMachine extends BaseStateMachine {
    // ...
	CompletableFuture<DataStream> stream(RaftClientRequest request) {
	  init = toProto(request);
	  dispatchCommand(init, streamInitContext());          // StreamInit, checks the token for init's block
	  channel = (KeyValueStreamDataChannel) dispatcher.getStreamDataChannel(init);   // no callback parameter
	  return new LocalStream(channel, getChunkExecutor(init), init,
	      dispatcher, getGroupId(), container2BCSIDMap);
	}
	
	CompletableFuture<?> link(DataStream stream, LogEntryProto entry) {
	  return ((LocalStream) stream).link();               // V3: fails unless the close PutBlock was committed
	}
class LocalStream implements StateMachine.DataStream {
  KeyValueStreamDataChannel channel;
  DataChannel dataChannel;           // what Ratis writes to and closes
  BlockID blockID;                   // the block the stream was opened for
  ContainerDispatcher dispatcher; RaftGroupId groupId; Map<Long, Long> container2BCSIDMap; Executor executor;

  LocalStream(channel, executor, init, dispatcher, groupId, container2BCSIDMap) {
    blockID = init.getWriteChunk().getBlockID();
    dataChannel = init.getCmdType() == StreamInitWithPutBlock
        ? new PutBlockOnClose()      // V3: commit the trailing PutBlock at close
        : channel;                   // V2: the channel drops the trailing PutBlock (HDDS-12007)
  }

  DataChannel getDataChannel() { return dataChannel; }

  // PutBlock in the middle of the stream (flush / hsync)
  CompletableFuture<ByteBuffer> onCommand(ByteBuffer command, long offset) {
    return supplyAsync(() -> {
      channel.drainBuffers();
      return putBlock(toProto(command)).toByteString().asReadOnlyByteBuffer();
    }, executor);
  }

  // The single place that commits a PutBlock received on the stream
  private ContainerCommandResponseProto putBlock(ContainerCommandRequestProto request) {
    require request.getCmdType() == PutBlock;
    require blockIdOf(request) == blockID;
    context = DispatcherContext.newBuilder(Op.STREAM_PUT_BLOCK)
        .setStage(COMBINED).setContainer2BCSIDMap(container2BCSIDMap).build();
    response = dispatcher.dispatch(request, context);
    if (response.getResult() != SUCCESS) {
      throw new StorageContainerException(response.getMessage(), response.getResult());
    }
    return response;
  }

  // Ratis closes the stream through getDataChannel().close()
  private class PutBlockOnClose implements DataChannel {
    write(b)  -> channel.write(b);
    force(m)  -> channel.force(m);
    isOpen()  -> channel.isOpen();
    close() {
      trailing = channel.closeAndReadPutBlock();   // writes the remaining data
      putBlock(trailing);                          // throws on failure, so the stream close fails
      channel.setLinked();
    }
  }

  CompletableFuture<?> link() { /* checks from ContainerStateMachine.link() */ }
  cleanUp(), getExecutor()   // unchanged
}

maybe in another PR.

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.

Sure we can keep this in mind and may refactor later.

case WRITE_STATE_MACHINE_DATA:
case READ_STATE_MACHINE_DATA:
case STREAM_LINK:
case STREAM_COMMAND:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Skipping this because sender has been verified first? so following message can be trusted unconditionaly?

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 is used for dispatch context but not the stream header:

    final DispatcherContext context = DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_COMMAND)
        .setStage(DispatcherContext.WriteChunkStage.COMBINED)
        .setContainer2BCSIDMap(container2BCSIDMap)
        .build();

…ozone/container/common/transport/server/ratis/LocalStream.java

Co-authored-by: Peter Lee <peterxcli@gmail.com>

@peterxcli peterxcli left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks! LGTM +1.

@amaliujia
amaliujia merged commit 7d9e841 into apache:master Oct 6, 2026
45 checks passed
errose28 added a commit that referenced this pull request Oct 6, 2026
* master: (64 commits)
  HDDS-16716. Add description for ozone.scm.ec.pipeline.per.volume.factor (#11415)
  HDDS-16008. PutBlocks from Flushes also go without Raft (#11356)
  HDDS-16666. Flush SCM transaction in memory during apply transaction (#11409)
  HDDS-15749. Run specific JUnit tests if possible (#10671)
  HDDS-16362. GetObjectAttributes ObjectParts should return Part entries for FSO buckets (#11242).
  HDDS-16643. Remove CleanupTableInfo mechanism (#11365)
  HDDS-16721. StreamBlockInputStream.read() returns a negative value for bytes 0x80 to 0xFF (#11411)
  HDDS-16241. gRPC deadline kills long-lived block streams after 30 seconds and the client never recovers (#11080)
  HDDS-15991. Speed up deleted table scans in quota repair (#11386)
  HDDS-16674. Bump awssdk to 2.55.6 (#11407)
  HDDS-16673. Avoid redundant ListBuckets RPCs when S3 bucket listing reaches the end (#11397)
  HDDS-16708. Let dependabot ignore iceberg minor version upgrades (#11398)
  HDDS-16713. Bump develocity-maven-extension to 2.6.0 (#11405)
  HDDS-16300. Allow OM to dynamically reconfigure its SCM node list without a restart (#11218)
  HDDS-15089. Support S3 per request read consistency (#11252)
  HDDS-16704. ReadBlock fails with IllegalStateException when a response is shorter than responseDataSize (#11402)
  HDDS-16631. Fix chooseRandom for rack names with common prefixes (#11401)
  HDDS-16654. Replace usage of deprecated finalize() in OM (#11376)
  HDDS-16658. Reuse source key details when opening input stream in S3 CopyObject (#11396)
  HDDS-16711. Bump moment to 2.31.0 (#11373)
  ...

Conflicts:
hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/HealthyReadOnlyNodeHandler.java
hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMStateMachine.java
hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDeadNodeHandler.java
hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestNodeStateManager.java
hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/rpc/TestRpcClient.java
hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/OmUtils.java
hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/ha/TestHadoopRpcOMFollowerReadFailoverProxyProvider.java
hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/shell/TestOzoneShellHA.java
hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/response/upgrade/OMCancelPrepareResponse.java
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/response/upgrade/OMCompleteFinalizeUpgradeResponse.java
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/response/upgrade/OMPrepareResponse.java
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java
hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants