mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Init proxyStartup.
This commit is contained in:
@@ -0,0 +1,56 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy;
|
||||
|
||||
import com.sun.net.httpserver.HttpExchange;
|
||||
import com.sun.net.httpserver.HttpHandler;
|
||||
import com.sun.net.httpserver.HttpServer;
|
||||
import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.net.InetSocketAddress;
|
||||
import org.apache.rocketmq.proxy.configuration.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.grpc.common.StartAndShutdown;
|
||||
|
||||
public class HealthCheckServer implements StartAndShutdown {
|
||||
|
||||
private HttpServer healthChecker;
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.healthChecker = HttpServer.create(new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0);
|
||||
this.healthChecker.createContext("/status", new HealthCheckHandler());
|
||||
this.healthChecker.setExecutor(null);
|
||||
this.healthChecker.start();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
this.healthChecker.stop(0);
|
||||
}
|
||||
|
||||
static class HealthCheckHandler implements HttpHandler {
|
||||
@Override
|
||||
public void handle(HttpExchange t) throws IOException {
|
||||
String response = "Hello";
|
||||
t.sendResponseHeaders(200, response.length());
|
||||
OutputStream os = t.getResponseBody();
|
||||
os.write(response.getBytes());
|
||||
os.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,125 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy;
|
||||
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.joran.JoranConfigurator;
|
||||
import ch.qos.logback.core.joran.spi.JoranException;
|
||||
import java.util.Date;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.broker.BrokerController;
|
||||
import org.apache.rocketmq.broker.BrokerStartup;
|
||||
import org.apache.rocketmq.client.log.ClientLogger;
|
||||
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
|
||||
import org.apache.rocketmq.proxy.configuration.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.configuration.ProxyConfig;
|
||||
import org.apache.rocketmq.proxy.grpc.GrpcServer;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.service.ClusterGrpcService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.GrpcService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.LocalGrpcService;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ProxyStartup {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ProxyStartup.class);
|
||||
|
||||
public static void main(String[] args) {
|
||||
try {
|
||||
ConfigurationManager.initEnv();
|
||||
initLogger();
|
||||
ConfigurationManager.intConfig();
|
||||
|
||||
// init thread pool monitor for proxy.
|
||||
initThreadPoolMonitor();
|
||||
|
||||
// create and start grpcServer
|
||||
GrpcServer grpcServer = createGrpcServer();
|
||||
grpcServer.start();
|
||||
|
||||
// health check server
|
||||
final HealthCheckServer healthCheckServer = new HealthCheckServer();
|
||||
healthCheckServer.start();
|
||||
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
|
||||
log.info("try to shutdown server");
|
||||
|
||||
try {
|
||||
healthCheckServer.shutdown();
|
||||
Thread.sleep(TimeUnit.SECONDS.toMillis(ConfigurationManager.getProxyConfig().getWaitAfterStopHealthCheckInSeconds()));
|
||||
} catch (Exception e) {
|
||||
log.error("err when shutdown healthCheckServer", e);
|
||||
}
|
||||
|
||||
try {
|
||||
grpcServer.shutdown();
|
||||
} catch (Exception e) {
|
||||
log.error("err when shutdown grpc server", e);
|
||||
}
|
||||
}));
|
||||
} catch (Exception e) {
|
||||
System.err.println("find a unexpect err." + e);
|
||||
e.printStackTrace();
|
||||
log.error("find a unexpect err.", e);
|
||||
System.exit(1);
|
||||
}
|
||||
|
||||
System.out.printf("%s%n", new Date() + " rmq-proxy startup successfully");
|
||||
log.info(new Date() + "rmq-proxy startup successfully");
|
||||
}
|
||||
|
||||
private static GrpcServer createGrpcServer() throws RuntimeException {
|
||||
GrpcService grpcService;
|
||||
String proxyModeStr = ConfigurationManager.getProxyConfig().getProxyMode();
|
||||
if (ProxyMode.isClusterMode(proxyModeStr)) {
|
||||
grpcService = new ClusterGrpcService();
|
||||
} else if (ProxyMode.isLocalMode(proxyModeStr)) {
|
||||
BrokerController brokerController = createBrokerController();
|
||||
grpcService = new LocalGrpcService(brokerController);
|
||||
} else {
|
||||
throw new IllegalArgumentException("try to start grpc server with wrong mode, use 'local' or 'cluster'");
|
||||
}
|
||||
|
||||
return new GrpcServer(grpcService);
|
||||
|
||||
}
|
||||
|
||||
private static BrokerController createBrokerController() {
|
||||
String[] brokerStartupArgs = new String[] {"-c", ConfigurationManager.getProxyConfig().getBrokerConfigPath()};
|
||||
return BrokerStartup.createBrokerController(brokerStartupArgs);
|
||||
}
|
||||
|
||||
private static void initThreadPoolMonitor() {
|
||||
ThreadPoolMonitor.init();
|
||||
ProxyConfig config = ConfigurationManager.getProxyConfig();
|
||||
ThreadPoolMonitor.config(config.isEnablePrintJstack(), config.getPrintJstackPeriodMillis());
|
||||
}
|
||||
|
||||
private static void initLogger() throws JoranException {
|
||||
System.setProperty(ClientLogger.CLIENT_LOG_USESLF4J, "true");
|
||||
|
||||
LoggerContext lc = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
JoranConfigurator configurator = new JoranConfigurator();
|
||||
configurator.setContext(lc);
|
||||
lc.reset();
|
||||
//https://logback.qos.ch/manual/configuration.html
|
||||
lc.setPackagingDataEnabled(false);
|
||||
configurator.doConfigure(ConfigurationManager.getProxyHome() + "/conf/logback.xml");
|
||||
}
|
||||
}
|
||||
@@ -22,24 +22,69 @@ import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
public class ProxyConfig {
|
||||
public final static String CONFIG_FILE_NAME = "rmq-proxy.json";
|
||||
|
||||
/**
|
||||
* Configuration for proxy
|
||||
*/
|
||||
private Integer healthCheckPort = 8000;
|
||||
private long waitAfterStopHealthCheckInSeconds = 40;
|
||||
|
||||
/**
|
||||
* configuration for ThreadPoolMonitor
|
||||
*/
|
||||
private boolean enablePrintJstack = true;
|
||||
private long printJstackPeriodMillis = 60000;
|
||||
|
||||
/**
|
||||
* gRPC
|
||||
*/
|
||||
private String proxyMode = ProxyMode.CLUSTER.name();
|
||||
private Boolean startGrpcServer = true;
|
||||
private Integer grpcServerPort = 8081;
|
||||
private String grpcTlsKeyPath = "/home/admin/rmq-gateway/conf/tls/gRPC.key.pem";
|
||||
private String grpcTlsCertPath = "/home/admin/rmq-gateway/conf/tls/gRPC.chain.cert.pem";
|
||||
private String grpcTlsKeyPath = ConfigurationManager.getProxyHome() + "/conf/tls/gRPC.key.pem";
|
||||
private String grpcTlsCertPath = ConfigurationManager.getProxyHome() + "/conf/tls/gRPC.chain.cert.pem";
|
||||
private int grpcBossLoopNum = 1;
|
||||
private int grpcWorkerLoopNum = Runtime.getRuntime().availableProcessors() * 2;
|
||||
private int grpcThreadPoolNums = 16 + Runtime.getRuntime().availableProcessors() * 2;
|
||||
private int grpcThreadPoolQueueCapacity = 100000;
|
||||
private String brokerConfigPath = ConfigurationManager.getProxyHome() + "/conf/broker.conf";
|
||||
/**
|
||||
* gRPC max message size
|
||||
* 130M = 4M * 32 messages + 2M attributes
|
||||
*/
|
||||
private int grpcMaxInboundMessageSize = 130 * 1024 * 1024;
|
||||
|
||||
public Integer getHealthCheckPort() {
|
||||
return healthCheckPort;
|
||||
}
|
||||
|
||||
public void setHealthCheckPort(Integer healthCheckPort) {
|
||||
this.healthCheckPort = healthCheckPort;
|
||||
}
|
||||
|
||||
public long getWaitAfterStopHealthCheckInSeconds() {
|
||||
return waitAfterStopHealthCheckInSeconds;
|
||||
}
|
||||
|
||||
public void setWaitAfterStopHealthCheckInSeconds(long waitAfterStopHealthCheckInSeconds) {
|
||||
this.waitAfterStopHealthCheckInSeconds = waitAfterStopHealthCheckInSeconds;
|
||||
}
|
||||
|
||||
public boolean isEnablePrintJstack() {
|
||||
return enablePrintJstack;
|
||||
}
|
||||
|
||||
public void setEnablePrintJstack(boolean enablePrintJstack) {
|
||||
this.enablePrintJstack = enablePrintJstack;
|
||||
}
|
||||
|
||||
public long getPrintJstackPeriodMillis() {
|
||||
return printJstackPeriodMillis;
|
||||
}
|
||||
|
||||
public void setPrintJstackPeriodMillis(long printJstackPeriodMillis) {
|
||||
this.printJstackPeriodMillis = printJstackPeriodMillis;
|
||||
}
|
||||
|
||||
public String getProxyMode() {
|
||||
return proxyMode;
|
||||
}
|
||||
@@ -112,6 +157,14 @@ public class ProxyConfig {
|
||||
this.grpcThreadPoolQueueCapacity = grpcThreadPoolQueueCapacity;
|
||||
}
|
||||
|
||||
public String getBrokerConfigPath() {
|
||||
return brokerConfigPath;
|
||||
}
|
||||
|
||||
public void setBrokerConfigPath(String brokerConfigPath) {
|
||||
this.brokerConfigPath = brokerConfigPath;
|
||||
}
|
||||
|
||||
public int getGrpcMaxInboundMessageSize() {
|
||||
return grpcMaxInboundMessageSize;
|
||||
}
|
||||
|
||||
@@ -42,8 +42,10 @@ public class GrpcServer {
|
||||
private final io.grpc.Server server;
|
||||
|
||||
private final ThreadPoolExecutor executor;
|
||||
private final GrpcService grpcService;
|
||||
|
||||
public GrpcServer(GrpcService grpcService) {
|
||||
this.grpcService = grpcService;
|
||||
int port = ConfigurationManager.getProxyConfig().getGrpcServerPort();
|
||||
NettyServerBuilder serverBuilder = NettyServerBuilder.forPort(port);
|
||||
|
||||
@@ -98,7 +100,11 @@ public class GrpcServer {
|
||||
bossLoopNum, workerLoopNum, maxInboundMessageSize);
|
||||
}
|
||||
|
||||
|
||||
public void start() throws Exception {
|
||||
// first to start grpc service.
|
||||
this.grpcService.start();
|
||||
|
||||
this.server.start();
|
||||
log.info("grpc server has started");
|
||||
}
|
||||
@@ -108,8 +114,10 @@ public class GrpcServer {
|
||||
this.server.shutdown().awaitTermination(30, TimeUnit.SECONDS);
|
||||
this.executor.shutdown();
|
||||
|
||||
this.grpcService.shutdown();
|
||||
|
||||
log.info("grpc server has stopped");
|
||||
} catch (InterruptedException e) {
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.common;
|
||||
|
||||
public interface StartAndShutdown {
|
||||
void start() throws Exception;
|
||||
|
||||
void shutdown() throws Exception;
|
||||
}
|
||||
@@ -64,19 +64,23 @@ public class ClusterGrpcService implements GrpcService {
|
||||
|
||||
}
|
||||
|
||||
@Override public CompleteFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
@Override
|
||||
public CompleteFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public CompleteFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
|
||||
@Override
|
||||
public CompleteFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public CompleteFuture<HealthCheckResponse> healthCheck(Context ctx, HealthCheckRequest request) {
|
||||
@Override
|
||||
public CompleteFuture<HealthCheckResponse> healthCheck(Context ctx, HealthCheckRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public CompleteFuture<SendMessageResponse> sendMessage(Context ctx, SendMessageRequest request) {
|
||||
@Override
|
||||
public CompleteFuture<SendMessageResponse> sendMessage(Context ctx, SendMessageRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -138,4 +142,12 @@ public class ClusterGrpcService implements GrpcService {
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() throws Exception {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,8 +53,9 @@ import apache.rocketmq.v1.SendMessageRequest;
|
||||
import apache.rocketmq.v1.SendMessageResponse;
|
||||
import io.grpc.Context;
|
||||
import io.netty.util.concurrent.CompleteFuture;
|
||||
import org.apache.rocketmq.proxy.grpc.common.StartAndShutdown;
|
||||
|
||||
public interface GrpcService {
|
||||
public interface GrpcService extends StartAndShutdown {
|
||||
CompleteFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request);
|
||||
|
||||
CompleteFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request);
|
||||
|
||||
@@ -140,4 +140,12 @@ public class LocalGrpcService implements GrpcService {
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public void start() throws Exception {
|
||||
this.brokerController.start();
|
||||
}
|
||||
|
||||
@Override public void shutdown() throws Exception {
|
||||
this.brokerController.shutdown();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user