szetszwo commented on code in PR #1534:
URL: https://github.com/apache/ratis/pull/1534#discussion_r3692194334
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -86,23 +87,65 @@
public class DataStreamManagement {
public static final Logger LOG =
LoggerFactory.getLogger(DataStreamManagement.class);
+ static final class LocalResult {
+ static final LocalResult ZERO = of(0);
+
+ static LocalResult of(long byteWritten) {
+ return new LocalResult(byteWritten, null);
+ }
+
+ static LocalResult of(ByteBuffer reply) {
+ return new LocalResult(0, reply);
+ }
+
+ /** For {@link Type#STREAM_DATA}. */
+ private final long byteWritten;
+ /** For {@link Type#STREAM_COMMAND}. */
+ private final ByteBuffer commandReply;
+
+ private LocalResult(long byteWritten, ByteBuffer commandReply) {
+ this.byteWritten = byteWritten;
+ this.commandReply = commandReply;
+ }
+
+ long getByteWritten() {
+ return byteWritten;
+ }
+
+ ByteBuffer getCommandReply() {
+ return commandReply;
+ }
+ }
+
static class LocalStream {
private final CompletableFuture<DataStream> streamFuture;
- private final AtomicReference<CompletableFuture<Long>> writeFuture;
- private final RequestMetrics metrics;
+ private final AtomicReference<CompletableFuture<LocalResult>> writeFuture;
+ private final RequestMetrics writeMetrics;
+ private final RequestMetrics commandMetrics;
- LocalStream(CompletableFuture<DataStream> streamFuture, RequestMetrics
metrics) {
+ LocalStream(CompletableFuture<DataStream> streamFuture, RequestMetrics
writeMetrics,
+ RequestMetrics commandMetrics) {
this.streamFuture = streamFuture;
- this.writeFuture = new AtomicReference<>(streamFuture.thenApply(s ->
0L));
- this.metrics = metrics;
+ this.writeFuture = new AtomicReference<>(streamFuture.thenApply(s ->
LocalResult.ZERO));
+ this.writeMetrics = writeMetrics;
+ this.commandMetrics = commandMetrics;
}
- CompletableFuture<Long> write(ByteBuf buf, Iterable<WriteOption> options,
+ CompletableFuture<LocalResult> write(ByteBuf buf, Iterable<WriteOption>
options,
Executor executor) {
Review Comment:
Make it a single line.
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -303,6 +359,22 @@ static CompletableFuture<Long> writeToAsync(ByteBuf buf,
return CompletableFuture.supplyAsync(() -> writeTo(buf, options, stream),
e);
}
+ static CompletableFuture<ByteBuffer> commandToAsync(ByteBuffer command, long
streamOffset,
+ DataStream stream, Executor defaultExecutor) {
+ final Executor e =
Optional.ofNullable(stream.getExecutor()).orElse(defaultExecutor);
+ return CompletableFuture.runAsync(() -> {}, e)
+ .thenCompose(ignored -> stream.onCommand(command, streamOffset));
+ }
+
+ static ByteBuffer copyBuffer(ByteBuf buf) {
+ final ByteBuffer copy = ByteBuffer.allocate(buf.readableBytes());
+ for (ByteBuffer buffer : buf.nioBuffers()) {
+ copy.put(buffer);
+ }
+ copy.flip();
+ return copy;
Review Comment:
Make it readonly
```java
return copy.asReadOnlyBuffer();
```
##########
ratis-client/src/main/java/org/apache/ratis/client/impl/DataStreamClientImpl.java:
##########
@@ -187,6 +199,11 @@ public CompletableFuture<DataStreamReply>
writeAsync(FilePositionCount src, Writ
return writeAsyncImpl(src, src.getCount(), Arrays.asList(options));
}
+ @Override
+ public CompletableFuture<DataStreamReply> commandAsync(ByteBuffer src) {
+ return commandAsyncImpl(src, src.remaining());
Review Comment:
commandAsyncImpl is used just once. Let's inline the code.
```java
public CompletableFuture<DataStreamReply> commandAsync(ByteBuffer src) {
if (isClosed()) {
return JavaUtils.completeExceptionally(new AlreadyClosedException(
clientId + ": stream already closed, request=" + header));
}
return combineHeader(send(Type.STREAM_COMMAND, src, src.remaining(),
Collections.singleton(StandardWriteOption.FLUSH)));
}
```
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -519,30 +601,48 @@ static void assertReplyCorrespondingToRequest(
Preconditions.assertTrue(request.getStreamOffset() ==
reply.getStreamOffset());
}
- private boolean
checkSuccessRemoteWrite(List<CompletableFuture<DataStreamReply>> replyFutures,
long bytesWritten,
- final DataStreamRequestByteBuf request) {
+ private boolean
checkSuccessRemoteWrite(List<CompletableFuture<DataStreamReply>> replyFutures,
+ LocalResult localResult, final DataStreamRequestByteBuf request) {
for (CompletableFuture<DataStreamReply> replyFuture : replyFutures) {
final DataStreamReply reply;
try {
reply = replyFuture.get(requestTimeout.getDuration(),
requestTimeout.getUnit());
} catch (Exception e) {
- throw new CompletionException("Failed to get reply for bytesWritten="
+ bytesWritten + ", " + request, e);
+ throw new CompletionException("Failed to get reply for " + localResult
+ ", " + request, e);
Review Comment:
Add toString()
```java
@Override
public String toString() {
return commandReply != null ? "commandReply:" +
StringUtils.bytes2HexString(commandReply)
: "byteWritten:" + byteWritten;
}
```
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -303,6 +359,22 @@ static CompletableFuture<Long> writeToAsync(ByteBuf buf,
return CompletableFuture.supplyAsync(() -> writeTo(buf, options, stream),
e);
}
+ static CompletableFuture<ByteBuffer> commandToAsync(ByteBuffer command, long
streamOffset,
+ DataStream stream, Executor defaultExecutor) {
+ final Executor e =
Optional.ofNullable(stream.getExecutor()).orElse(defaultExecutor);
+ return CompletableFuture.runAsync(() -> {}, e)
+ .thenCompose(ignored -> stream.onCommand(command, streamOffset));
Review Comment:
Using completedFuture(null) is more efficient since it doesn't have to run
an empty task.
```java
return CompletableFuture.completedFuture(null)
.thenComposeAsync(ignored -> stream.onCommand(command,
streamOffset), e);
```
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -469,24 +544,31 @@ private void readImpl(DataStreamRequestByteBuf request,
ChannelHandlerContext ct
() -> new IllegalStateException("Failed to get StreamInfo for " +
request));
}
- final CompletableFuture<Long> localWrite;
+ final CompletableFuture<LocalResult> localResult;
final List<CompletableFuture<DataStreamReply>> remoteWrites;
if (request.getType() == Type.STREAM_HEADER) {
- localWrite = CompletableFuture.completedFuture(0L);
+ localResult = CompletableFuture.completedFuture(LocalResult.ZERO);
remoteWrites = Collections.emptyList();
} else if (request.getType() == Type.STREAM_DATA) {
- localWrite = info.getLocal().write(request.slice(),
request.getWriteOptionList(), writeExecutor);
+ localResult = info.getLocal().write(request.slice(),
request.getWriteOptionList(), writeExecutor);
remoteWrites = info.applyToRemotes(out -> out.write(request,
requestExecutor));
+ } else if (request.getType() == Type.STREAM_COMMAND) {
+ // command is supposed to have small data size, just copy it
+ final ByteBuffer command = copyBuffer(request.slice());
+ localResult = info.getLocal().command(command,
request.getStreamOffset(), writeExecutor);
+ remoteWrites = info.applyToRemotes(out -> out.command(
+ copyBuffer(request.slice()), requestExecutor));
Review Comment:
duplicate() and reuse command, i.e.
```java
final ByteBuffer command = copyBuffer(request.slice());
localResult = info.getLocal().command(command.duplicate(),
request.getStreamOffset(), writeExecutor);
remoteWrites = info.applyToRemotes(out -> out.command(command,
requestExecutor));
```
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -86,23 +87,65 @@
public class DataStreamManagement {
public static final Logger LOG =
LoggerFactory.getLogger(DataStreamManagement.class);
+ static final class LocalResult {
+ static final LocalResult ZERO = of(0);
+
+ static LocalResult of(long byteWritten) {
+ return new LocalResult(byteWritten, null);
+ }
+
+ static LocalResult of(ByteBuffer reply) {
+ return new LocalResult(0, reply);
+ }
+
+ /** For {@link Type#STREAM_DATA}. */
+ private final long byteWritten;
+ /** For {@link Type#STREAM_COMMAND}. */
+ private final ByteBuffer commandReply;
+
+ private LocalResult(long byteWritten, ByteBuffer commandReply) {
+ this.byteWritten = byteWritten;
+ this.commandReply = commandReply;
+ }
+
+ long getByteWritten() {
+ return byteWritten;
+ }
+
+ ByteBuffer getCommandReply() {
+ return commandReply;
+ }
+ }
+
static class LocalStream {
private final CompletableFuture<DataStream> streamFuture;
- private final AtomicReference<CompletableFuture<Long>> writeFuture;
- private final RequestMetrics metrics;
+ private final AtomicReference<CompletableFuture<LocalResult>> writeFuture;
Review Comment:
Let's rename it to `resultFuture`.
--
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]