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;
   }
 

Reply via email to