szetszwo commented on code in PR #11356:
URL: https://github.com/apache/ozone/pull/11356#discussion_r4168207337


##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java:
##########
@@ -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);

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



##########
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",

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



##########
hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java:
##########
@@ -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());

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



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to