Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions ratis-netty/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,11 @@
<artifactId>junit-platform-launcher</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</dependency>

</dependencies>

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,29 @@ static boolean useEpoll(RaftProperties properties) {
static void setUseEpoll(RaftProperties properties, boolean enable) {
setBoolean(properties::setBoolean, USE_EPOLL_KEY, enable);
}

String ASYNC_REQUEST_THREAD_POOL_CACHED_KEY = PREFIX + ".async.request.thread.pool.cached";
// Default to a fixed pool.
// TODO: Refer to https://issues.apache.org/jira/browse/RATIS-2637
boolean ASYNC_REQUEST_THREAD_POOL_CACHED_DEFAULT = false;
static boolean asyncRequestThreadPoolCached(RaftProperties properties) {
return getBoolean(properties::getBoolean, ASYNC_REQUEST_THREAD_POOL_CACHED_KEY,
ASYNC_REQUEST_THREAD_POOL_CACHED_DEFAULT, getDefaultLog());
}
static void setAsyncRequestThreadPoolCached(RaftProperties properties, boolean useCached) {
setBoolean(properties::setBoolean, ASYNC_REQUEST_THREAD_POOL_CACHED_KEY, useCached);
}

String ASYNC_REQUEST_THREAD_POOL_SIZE_KEY = PREFIX + ".async.request.thread.pool.size";
int ASYNC_REQUEST_THREAD_POOL_SIZE_DEFAULT = 32;
static int asyncRequestThreadPoolSize(RaftProperties properties) {
return getInt(properties::getInt, ASYNC_REQUEST_THREAD_POOL_SIZE_KEY,
ASYNC_REQUEST_THREAD_POOL_SIZE_DEFAULT, getDefaultLog(),
requireMin(0), requireMax(65536));
}
static void setAsyncRequestThreadPoolSize(RaftProperties properties, int size) {
setInt(properties::setInt, ASYNC_REQUEST_THREAD_POOL_SIZE_KEY, size);
}
}

interface Client {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,6 @@
import org.apache.ratis.netty.NettyConfigKeys;
import org.apache.ratis.netty.NettyRpcProxy;
import org.apache.ratis.util.NettyUtils;
import org.apache.ratis.protocol.GroupInfoReply;
import org.apache.ratis.protocol.GroupListReply;
import org.apache.ratis.protocol.RaftClientReply;
import org.apache.ratis.protocol.RaftPeerId;
import org.apache.ratis.rpc.SupportedRpcType;
Expand All @@ -42,6 +40,7 @@
import org.apache.ratis.proto.netty.NettyProtos.RaftNettyServerReplyProto;
import org.apache.ratis.proto.netty.NettyProtos.RaftNettyServerRequestProto;
import org.apache.ratis.util.CodeInjectionForTesting;
import org.apache.ratis.util.ConcurrentUtils;
import org.apache.ratis.util.JavaUtils;
import org.apache.ratis.util.MemoizedSupplier;
import org.apache.ratis.util.ProtoUtils;
Expand All @@ -50,8 +49,10 @@

import java.io.IOException;
import java.net.InetSocketAddress;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.Objects;

/**
* A netty server endpoint that acts as the communication layer.
Expand Down Expand Up @@ -87,12 +88,43 @@ public static Builder newBuilder() {
private final MemoizedSupplier<ChannelFuture> channel;
private final InetSocketAddress socketAddress;

@ChannelHandler.Sharable
private final ExecutorService requestExecutor;

class InboundHandler extends SimpleChannelInboundHandler<RaftNettyServerRequestProto> {
/**
* Tail of this channel's chain of request-handling tasks.
* Requests on a channel must be handled in arrival order.
*/
private CompletableFuture<Void> tail = CompletableFuture.completedFuture(null);

@Override
protected void channelRead0(ChannelHandlerContext ctx, RaftNettyServerRequestProto proto) {
final RaftNettyServerReplyProto reply = handle(proto);
ctx.writeAndFlush(reply);
tail = tail.handleAsync((prev, prevError) -> {
final CompletableFuture<RaftNettyServerReplyProto> replyFuture;
try {
replyFuture = handleAsync(proto);
} catch (Exception e) {
// No request context to build a reply; fail fast by closing the channel.
LOG.warn("{}: Failed to handle request; closing the channel.", getId(), e);
ctx.close();
return null;
}
replyFuture.whenComplete((reply, e) -> {
if (e != null) {
LOG.warn("{}: Failed to handle request; closing the channel.", getId(), e);
ctx.close();
} else {
ctx.writeAndFlush(reply);
}
});
return null;
}, requestExecutor);
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
LOG.warn("{}: exceptionCaught on channel {}; closing it.", getId(), ctx.channel(), cause);
ctx.close();
}
}

Expand All @@ -116,6 +148,11 @@ protected void initChannel(SocketChannel ch) {
}
};

this.requestExecutor = ConcurrentUtils.newThreadPoolWithMax(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

So, it's still a single worker pool (corePoolSize is 0). But instead of using TCP backpressure, we are pushing everything to the unlimited queue and creating pressure on the memory. Moreover, previously, Netty’s worker EventLoops handled connection shards independently, so a blocked request delayed only that shard’s RPCs. The new default single request worker queues all inbound RPCs, including heartbeats, so one slow request can delay heartbeats and trigger leader election.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes this was one issue which I missed and faced before (hence the test failure in flaky test suite).
corePoolSize=0 + an unbounded queue causes requestExecutor to be single-threaded, and since handle() blocks until commit, this causes the timeouts.
I have addressed this by switching to a fixed pool for now as the related change would increase LoC.

Filed https://issues.apache.org/jira/browse/RATIS-2637 for the improvement as a follow up.

NettyConfigKeys.Server.asyncRequestThreadPoolCached(server.getProperties()),
NettyConfigKeys.Server.asyncRequestThreadPoolSize(server.getProperties()),
server.getId() + "-request-");

final boolean useEpoll = NettyConfigKeys.Server.useEpoll(server.getProperties());
this.bossGroup = NettyUtils.newEventLoopGroup(CLASS_NAME + "-bossGroup", 0, useEpoll);
this.workerGroup = NettyUtils.newEventLoopGroup(CLASS_NAME + "-workerGroup",0, useEpoll);
Expand Down Expand Up @@ -155,6 +192,7 @@ public void startImpl() throws IOException {

@Override
public void closeImpl() throws IOException {
ConcurrentUtils.shutdownAndWait(requestExecutor);
final ChannelFuture f = getChannel().close();
f.syncUninterruptibly();
bossGroup.shutdownGracefully(0, 100, TimeUnit.MILLISECONDS);
Expand All @@ -181,113 +219,111 @@ public InetSocketAddress getInetSocketAddress() {
}
}

RaftNettyServerReplyProto handle(RaftNettyServerRequestProto proto) {
CompletableFuture<RaftNettyServerReplyProto> handleAsync(RaftNettyServerRequestProto proto) {
RaftRpcRequestProto rpcRequest = null;
try {
final CompletableFuture<RaftNettyServerReplyProto> replyFuture;
switch (proto.getRaftNettyServerRequestCase()) {
case REQUESTVOTEREQUEST:
final RequestVoteRequestProto request = proto.getRequestVoteRequest();
rpcRequest = request.getServerRequest();
final RequestVoteReplyProto reply = server.requestVote(request);
return RaftNettyServerReplyProto.newBuilder()
.setRequestVoteReply(reply)
.build();
// requestVote has no async variant; it is fast and does not block on commit.
final RequestVoteRequestProto requestVoteRequest = proto.getRequestVoteRequest();
rpcRequest = requestVoteRequest.getServerRequest();
replyFuture = CompletableFuture.completedFuture(RaftNettyServerReplyProto.newBuilder()
.setRequestVoteReply(server.requestVote(requestVoteRequest))
.build());
break;

case TRANSFERLEADERSHIPREQUEST:
final TransferLeadershipRequestProto transferLeadershipRequest = proto.getTransferLeadershipRequest();
rpcRequest = transferLeadershipRequest.getRpcRequest();
final RaftClientReply transferLeadershipReply = server.transferLeadership(
ClientProtoUtils.toTransferLeadershipRequest(transferLeadershipRequest));
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(transferLeadershipReply))
.build();
replyFuture = server.transferLeadershipAsync(
ClientProtoUtils.toTransferLeadershipRequest(transferLeadershipRequest))
.thenApply(NettyRpcService::toRaftClientReply);
break;

case STARTLEADERELECTIONREQUEST:
// startLeaderElection has no async variant; it is fast and does not block on commit.
final StartLeaderElectionRequestProto startLeaderElectionRequest = proto.getStartLeaderElectionRequest();
rpcRequest = startLeaderElectionRequest.getServerRequest();
final StartLeaderElectionReplyProto startLeaderElectionReply =
server.startLeaderElection(startLeaderElectionRequest);
return RaftNettyServerReplyProto.newBuilder().setStartLeaderElectionReply(startLeaderElectionReply).build();
replyFuture = CompletableFuture.completedFuture(RaftNettyServerReplyProto.newBuilder()
.setStartLeaderElectionReply(server.startLeaderElection(startLeaderElectionRequest))
.build());
break;

case SNAPSHOTMANAGEMENTREQUEST:
final SnapshotManagementRequestProto snapshotManagementRequest = proto.getSnapshotManagementRequest();
rpcRequest = snapshotManagementRequest.getRpcRequest();
final RaftClientReply snapshotManagementReply = server.snapshotManagement(
ClientProtoUtils.toSnapshotManagementRequest(snapshotManagementRequest));
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(snapshotManagementReply))
.build();
replyFuture = server.snapshotManagementAsync(
ClientProtoUtils.toSnapshotManagementRequest(snapshotManagementRequest))
.thenApply(NettyRpcService::toRaftClientReply);
break;

case LEADERELECTIONMANAGEMENTREQUEST:
final LeaderElectionManagementRequestProto leaderElectionManagementRequest =
proto.getLeaderElectionManagementRequest();
rpcRequest = leaderElectionManagementRequest.getRpcRequest();
final RaftClientReply leaderElectionManagementReply = server.leaderElectionManagement(
ClientProtoUtils.toLeaderElectionManagementRequest(leaderElectionManagementRequest));
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(leaderElectionManagementReply))
.build();
replyFuture = server.leaderElectionManagementAsync(
ClientProtoUtils.toLeaderElectionManagementRequest(leaderElectionManagementRequest))
.thenApply(NettyRpcService::toRaftClientReply);
break;

case APPENDENTRIESREQUEST:
final AppendEntriesRequestProto appendEntriesRequest = proto.getAppendEntriesRequest();
rpcRequest = appendEntriesRequest.getServerRequest();
final AppendEntriesReplyProto appendEntriesReply = server.appendEntries(appendEntriesRequest);
return RaftNettyServerReplyProto.newBuilder()
.setAppendEntriesReply(appendEntriesReply)
.build();
replyFuture = server.appendEntriesAsync(appendEntriesRequest)
.thenApply(reply -> RaftNettyServerReplyProto.newBuilder()
.setAppendEntriesReply(reply)
.build());
break;

case INSTALLSNAPSHOTREQUEST:
// installSnapshot has no async variant; it runs on this per-channel serialized path.
final InstallSnapshotRequestProto installSnapshotRequest = proto.getInstallSnapshotRequest();
rpcRequest = installSnapshotRequest.getServerRequest();
final InstallSnapshotReplyProto installSnapshotReply = server.installSnapshot(installSnapshotRequest);
return RaftNettyServerReplyProto.newBuilder()
.setInstallSnapshotReply(installSnapshotReply)
.build();
replyFuture = CompletableFuture.completedFuture(RaftNettyServerReplyProto.newBuilder()
.setInstallSnapshotReply(server.installSnapshot(installSnapshotRequest))
.build());
break;

case RAFTCLIENTREQUEST:
final RaftClientRequestProto raftClientRequest = proto.getRaftClientRequest();
rpcRequest = raftClientRequest.getRpcRequest();
final RaftClientReply raftClientReply = server.submitClientRequest(
ClientProtoUtils.toRaftClientRequest(raftClientRequest));
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(raftClientReply))
.build();
replyFuture = server.submitClientRequestAsync(ClientProtoUtils.toRaftClientRequest(raftClientRequest))
.thenApply(NettyRpcService::toRaftClientReply);
break;

case SETCONFIGURATIONREQUEST:
final SetConfigurationRequestProto configurationRequest = proto.getSetConfigurationRequest();
rpcRequest = configurationRequest.getRpcRequest();
final RaftClientReply configurationReply = server.setConfiguration(
ClientProtoUtils.toSetConfigurationRequest(configurationRequest));
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(configurationReply))
.build();
final SetConfigurationRequestProto setConfigurationRequest = proto.getSetConfigurationRequest();
rpcRequest = setConfigurationRequest.getRpcRequest();
replyFuture = server.setConfigurationAsync(
ClientProtoUtils.toSetConfigurationRequest(setConfigurationRequest))
.thenApply(NettyRpcService::toRaftClientReply);
break;

case GROUPMANAGEMENTREQUEST:
final GroupManagementRequestProto groupManagementRequest = proto.getGroupManagementRequest();
rpcRequest = groupManagementRequest.getRpcRequest();
final RaftClientReply groupManagementReply = server.groupManagement(
ClientProtoUtils.toGroupManagementRequest(groupManagementRequest));
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(groupManagementReply))
.build();
replyFuture = server.groupManagementAsync(ClientProtoUtils.toGroupManagementRequest(groupManagementRequest))
.thenApply(NettyRpcService::toRaftClientReply);
break;

case GROUPLISTREQUEST:
final GroupListRequestProto groupListRequest = proto.getGroupListRequest();
rpcRequest = groupListRequest.getRpcRequest();
final GroupListReply groupListReply = server.getGroupList(
ClientProtoUtils.toGroupListRequest(groupListRequest));
return RaftNettyServerReplyProto.newBuilder()
.setGroupListReply(ClientProtoUtils.toGroupListReplyProto(groupListReply))
.build();
replyFuture = server.getGroupListAsync(ClientProtoUtils.toGroupListRequest(groupListRequest))
.thenApply(reply -> RaftNettyServerReplyProto.newBuilder()
.setGroupListReply(ClientProtoUtils.toGroupListReplyProto(reply))
.build());
break;

case GROUPINFOREQUEST:
final GroupInfoRequestProto groupInfoRequest = proto.getGroupInfoRequest();
rpcRequest = groupInfoRequest.getRpcRequest();
final GroupInfoReply groupInfoReply = server.getGroupInfo(
ClientProtoUtils.toGroupInfoRequest(groupInfoRequest));
return RaftNettyServerReplyProto.newBuilder()
.setGroupInfoReply(ClientProtoUtils.toGroupInfoReplyProto(groupInfoReply))
.build();
replyFuture = server.getGroupInfoAsync(ClientProtoUtils.toGroupInfoRequest(groupInfoRequest))
.thenApply(reply -> RaftNettyServerReplyProto.newBuilder()
.setGroupInfoReply(ClientProtoUtils.toGroupInfoReplyProto(reply))
.build());
break;

case RAFTNETTYSERVERREQUEST_NOT_SET:
throw new IllegalArgumentException("Request case not set in proto: "
Expand All @@ -296,12 +332,31 @@ RaftNettyServerReplyProto handle(RaftNettyServerRequestProto proto) {
throw new UnsupportedOperationException("Request case not supported: "
+ proto.getRaftNettyServerRequestCase());
}
} catch (IOException ioe) {
return toRaftNettyServerReplyProto(
Objects.requireNonNull(rpcRequest, "rpcRequest = null"), ioe);

final RaftRpcRequestProto request = rpcRequest;
// Convert an asynchronous failure into an error reply (the client casts it to IOException).
return replyFuture.exceptionally(e -> toRaftNettyServerReplyProto(request, toIOException(e)));
} catch (IOException | RuntimeException e) {
// A synchronous failure before the reply future was created.
if (rpcRequest == null) {
// No request context to build a targeted reply; let InboundHandler close the channel.
throw new IllegalStateException(getId() + ": Failed to handle request " + proto, e);
}
return CompletableFuture.completedFuture(toRaftNettyServerReplyProto(rpcRequest, toIOException(e)));
}
}

private static RaftNettyServerReplyProto toRaftClientReply(RaftClientReply reply) {
return RaftNettyServerReplyProto.newBuilder()
.setRaftClientReply(ClientProtoUtils.toRaftClientReplyProto(reply))
.build();
}

private static IOException toIOException(Throwable t) {
final Throwable cause = t instanceof CompletionException && t.getCause() != null ? t.getCause() : t;
return cause instanceof IOException ? (IOException) cause : new IOException(cause);
}

private static RaftNettyServerReplyProto toRaftNettyServerReplyProto(
RaftRpcRequestProto request, IOException e) {
final RaftRpcReplyProto.Builder rpcReply = RaftRpcReplyProto.newBuilder()
Expand Down
Loading
Loading