Repository: incubator-ratis Updated Branches: refs/heads/master 415bb3ed7 -> bb5c258c3
RATIS-197. In gRPC, reads are failing for size greater than 4MB. Contributed by Mukul Kumar Singh Project: http://git-wip-us.apache.org/repos/asf/incubator-ratis/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-ratis/commit/bb5c258c Tree: http://git-wip-us.apache.org/repos/asf/incubator-ratis/tree/bb5c258c Diff: http://git-wip-us.apache.org/repos/asf/incubator-ratis/diff/bb5c258c Branch: refs/heads/master Commit: bb5c258c308acb04bf0daf412f89a3da8400f694 Parents: 415bb3e Author: Tsz-Wo Nicholas Sze <[email protected]> Authored: Wed Feb 7 12:20:28 2018 -0800 Committer: Tsz-Wo Nicholas Sze <[email protected]> Committed: Wed Feb 7 12:20:28 2018 -0800 ---------------------------------------------------------------------- .../java/org/apache/ratis/grpc/client/AppendStreamer.java | 10 +++++----- .../java/org/apache/ratis/grpc/client/GrpcClientRpc.java | 6 +++--- .../ratis/grpc/client/RaftClientProtocolClient.java | 3 ++- .../apache/ratis/grpc/client/RaftClientProtocolProxy.java | 4 ++-- 4 files changed, 12 insertions(+), 11 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/bb5c258c/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/AppendStreamer.java ---------------------------------------------------------------------- diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/AppendStreamer.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/AppendStreamer.java index 281c838..08e8376 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/AppendStreamer.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/AppendStreamer.java @@ -71,7 +71,7 @@ public class AppendStreamer implements Closeable { private final Deque<RaftClientRequestProto> dataQueue; private final Deque<RaftClientRequestProto> ackQueue; private final int maxPendingNum; - private final int maxMessageSize; + private final SizeInBytes maxMessageSize; private final PeerProxyMap<RaftClientProtocolProxy> proxyMap; private final Map<RaftPeerId, RaftPeer> peers; @@ -88,7 +88,7 @@ public class AppendStreamer implements Closeable { RaftPeerId leaderId, ClientId clientId) { this.clientId = clientId; maxPendingNum = GrpcConfigKeys.OutputStream.outstandingAppendsMax(prop); - maxMessageSize = GrpcConfigKeys.messageSizeMax(prop).getSizeInt(); + maxMessageSize = GrpcConfigKeys.messageSizeMax(prop); dataQueue = new ConcurrentLinkedDeque<>(); ackQueue = new ConcurrentLinkedDeque<>(); exceptionAndRetry = new ExceptionAndRetry(prop); @@ -98,7 +98,7 @@ public class AppendStreamer implements Closeable { Collectors.toMap(RaftPeer::getId, Function.identity())); proxyMap = new PeerProxyMap<>(clientId.toString(), raftPeer -> new RaftClientProtocolProxy(clientId, raftPeer, ResponseHandler::new, - GrpcConfigKeys.flowControlWindow(prop))); + GrpcConfigKeys.flowControlWindow(prop), maxMessageSize)); proxyMap.addPeers(group.getPeers()); refreshLeaderProxy(leaderId, null); @@ -158,9 +158,9 @@ public class AppendStreamer implements Closeable { // wrap the current buffer into a RaftClientRequestProto final RaftClientRequestProto request = ClientProtoUtils.toRaftClientRequestProto( clientId, leaderId, groupId, seqNum, seqNum, content, false); - if (request.getSerializedSize() > maxMessageSize) { + if (request.getSerializedSize() > maxMessageSize.getSizeInt()) { throw new IOException("msg size:" + request.getSerializedSize() + - " exceeds maximum:" + maxMessageSize); + " exceeds maximum:" + maxMessageSize.getSizeInt()); } dataQueue.offer(request); this.notifyAll(); http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/bb5c258c/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/GrpcClientRpc.java ---------------------------------------------------------------------- diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/GrpcClientRpc.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/GrpcClientRpc.java index 508260e..c235997 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/GrpcClientRpc.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/GrpcClientRpc.java @@ -49,7 +49,7 @@ public class GrpcClientRpc extends RaftClientRpcWithProxy<RaftClientProtocolClie public GrpcClientRpc(ClientId clientId, RaftProperties properties) { super(new PeerProxyMap<>(clientId.toString(), p -> new RaftClientProtocolClient(clientId, p, - GrpcConfigKeys.flowControlWindow(properties)))); + GrpcConfigKeys.flowControlWindow(properties), GrpcConfigKeys.messageSizeMax(properties)))); this.clientId = clientId; maxMessageSize = GrpcConfigKeys.messageSizeMax(properties).getSizeInt(); } @@ -129,10 +129,10 @@ public class GrpcClientRpc extends RaftClientRpcWithProxy<RaftClientProtocolClie requestObserver.onNext(requestProto); requestObserver.onCompleted(); - return replyFuture.thenApply(replyProto -> toRaftClientReply(replyProto)); + return replyFuture.thenApply(ClientProtoUtils::toRaftClientReply); } - RaftClientRequestProto toRaftClientRequestProto(RaftClientRequest request) throws IOException { + private RaftClientRequestProto toRaftClientRequestProto(RaftClientRequest request) throws IOException { final RaftClientRequestProto proto = ClientProtoUtils.toRaftClientRequestProto(request); if (proto.getSerializedSize() > maxMessageSize) { throw new IOException(clientId + ": Message size:" + proto.getSerializedSize() http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/bb5c258c/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java ---------------------------------------------------------------------- diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java index 601232a..97fc9a3 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java @@ -58,11 +58,12 @@ public class RaftClientProtocolClient implements Closeable { private final AtomicReference<AsyncStreamObservers> appendStreamObservers = new AtomicReference<>(); public RaftClientProtocolClient(ClientId id, RaftPeer target, - SizeInBytes flowControlWindow) { + SizeInBytes flowControlWindow, SizeInBytes maxMessageSize) { this.name = JavaUtils.memoize(() -> id + "->" + target.getId()); this.target = target; channel = NettyChannelBuilder.forTarget(target.getAddress()) .usePlaintext(true).flowControlWindow(flowControlWindow.getSizeInt()) + .maxMessageSize(maxMessageSize.getSizeInt()) .build(); blockingStub = RaftClientProtocolServiceGrpc.newBlockingStub(channel); asyncStub = RaftClientProtocolServiceGrpc.newStub(channel); http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/bb5c258c/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolProxy.java ---------------------------------------------------------------------- diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolProxy.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolProxy.java index d07add7..2ed70c2 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolProxy.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolProxy.java @@ -36,8 +36,8 @@ public class RaftClientProtocolProxy implements Closeable { public RaftClientProtocolProxy( ClientId clientId, RaftPeer target, Function<RaftPeer, CloseableStreamObserver> responseHandlerCreation, - SizeInBytes flowControlWindow) { - proxy = new RaftClientProtocolClient(clientId, target, flowControlWindow); + SizeInBytes flowControlWindow, SizeInBytes maxMessageSize) { + proxy = new RaftClientProtocolClient(clientId, target, flowControlWindow, maxMessageSize); this.responseHandlerCreation = responseHandlerCreation; }
