amaliujia commented on code in PR #11356:
URL: https://github.com/apache/ozone/pull/11356#discussion_r4180579516
##########
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:
done. though I need to add a lock to `byteBufferList` in this case.
##########
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:
done
--
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]