Repository navigation
HDDS-16008. PutBlocks from Flushes also go without Raft - #11356
Conversation
|
The core changes are minors. Majority LOC is about testing. |
szetszwo
left a comment
There was a problem hiding this comment.
@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", |
There was a problem hiding this comment.
We should add a new conf but not changing an existing conf. Otherwise, it becomes an incompatible change.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Oh, you are right! Then, we should change it instead.
| // 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); |
There was a problem hiding this comment.
Let's just release the buffer here without using commitWatcher?
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Then, let's use commitWatcher for now. We may simplify the code later on. Sorry that I was aware of the complication.
There was a problem hiding this comment.
Sure revert back to use commitWatcher
| final ContainerCommandRequestProto request = ContainerCommandRequestMessage.toProto( | ||
| ByteString.copyFrom(command), getGroupId()); |
There was a problem hiding this comment.
Add a new ContainerCommandRequestMessage.toProto(ByteBuffer, RaftGroupId) method to avoid copying.
There was a problem hiding this comment.
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());|
@amaliujia , thanks for the update! Please see the responses above. |
2df1f8b to
0eb6312
Compare
szetszwo
left a comment
There was a problem hiding this comment.
+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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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: |
There was a problem hiding this comment.
Skipping this because sender has been verified first? so following message can be trusted unconditionaly?
There was a problem hiding this comment.
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>
* 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
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