mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-19 02:23:24 +08:00
* Add timeout configuration for grpc server * Add proxyConfig
This commit is contained in:
@@ -85,6 +85,7 @@ public class ProxyStartup {
|
||||
.addService(ChannelzService.newInstance(100))
|
||||
.addService(ProtoReflectionService.newInstance())
|
||||
.configInterceptor(accessValidators)
|
||||
.shutdownTime(ConfigurationManager.getProxyConfig().getGrpcShutdownTimeSeconds(), TimeUnit.SECONDS)
|
||||
.build();
|
||||
PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(grpcServer);
|
||||
|
||||
|
||||
@@ -87,6 +87,7 @@ public class ProxyConfig implements ConfigFile {
|
||||
*/
|
||||
private String proxyMode = ProxyMode.CLUSTER.name();
|
||||
private Integer grpcServerPort = 8081;
|
||||
private long grpcShutdownTimeSeconds = 30;
|
||||
private int grpcBossLoopNum = 1;
|
||||
private int grpcWorkerLoopNum = PROCESSOR_NUMBER * 2;
|
||||
private boolean enableGrpcEpoll = false;
|
||||
@@ -443,6 +444,14 @@ public class ProxyConfig implements ConfigFile {
|
||||
this.grpcServerPort = grpcServerPort;
|
||||
}
|
||||
|
||||
public long getGrpcShutdownTimeSeconds() {
|
||||
return grpcShutdownTimeSeconds;
|
||||
}
|
||||
|
||||
public void setGrpcShutdownTimeSeconds(long grpcShutdownTimeSeconds) {
|
||||
this.grpcShutdownTimeSeconds = grpcShutdownTimeSeconds;
|
||||
}
|
||||
|
||||
public boolean isUseEndpointPortFromRequest() {
|
||||
return useEndpointPortFromRequest;
|
||||
}
|
||||
|
||||
@@ -29,8 +29,14 @@ public class GrpcServer implements StartAndShutdown {
|
||||
|
||||
private final Server server;
|
||||
|
||||
protected GrpcServer(Server server) {
|
||||
private final long timeout;
|
||||
|
||||
private final TimeUnit unit;
|
||||
|
||||
protected GrpcServer(Server server, long timeout, TimeUnit unit) {
|
||||
this.server = server;
|
||||
this.timeout = timeout;
|
||||
this.unit = unit;
|
||||
}
|
||||
|
||||
public void start() throws Exception {
|
||||
@@ -40,7 +46,7 @@ public class GrpcServer implements StartAndShutdown {
|
||||
|
||||
public void shutdown() {
|
||||
try {
|
||||
this.server.shutdown().awaitTermination(30, TimeUnit.SECONDS);
|
||||
this.server.shutdown().awaitTermination(timeout, unit);
|
||||
log.info("grpc server shutdown successfully.");
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
|
||||
@@ -41,6 +41,10 @@ public class GrpcServerBuilder {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
protected NettyServerBuilder serverBuilder;
|
||||
|
||||
protected long time = 30;
|
||||
|
||||
protected TimeUnit unit = TimeUnit.SECONDS;
|
||||
|
||||
public static GrpcServerBuilder newBuilder(ThreadPoolExecutor executor, int port) {
|
||||
return new GrpcServerBuilder(executor, port);
|
||||
}
|
||||
@@ -77,6 +81,12 @@ public class GrpcServerBuilder {
|
||||
port, bossLoopNum, workerLoopNum, maxInboundMessageSize);
|
||||
}
|
||||
|
||||
public GrpcServerBuilder shutdownTime(long time, TimeUnit unit) {
|
||||
this.time = time;
|
||||
this.unit = unit;
|
||||
return this;
|
||||
}
|
||||
|
||||
public GrpcServerBuilder addService(BindableService service) {
|
||||
this.serverBuilder.addService(service);
|
||||
return this;
|
||||
@@ -93,7 +103,7 @@ public class GrpcServerBuilder {
|
||||
}
|
||||
|
||||
public GrpcServer build() {
|
||||
return new GrpcServer(this.serverBuilder.build());
|
||||
return new GrpcServer(this.serverBuilder.build(), time, unit);
|
||||
}
|
||||
|
||||
public GrpcServerBuilder configInterceptor(List<AccessValidator> accessValidators) {
|
||||
|
||||
Reference in New Issue
Block a user