mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
@@ -45,5 +45,4 @@ public class LoggerName {
|
||||
public static final String FAILOVER_LOGGER_NAME = "RocketmqFailover";
|
||||
public static final String STDOUT_LOGGER_NAME = "STDOUT";
|
||||
public static final String PROXY_LOGGER_NAME = "RocketmqProxy";
|
||||
public static final String GRPC_LOGGER_NAME = "RocketmqGrpc";
|
||||
}
|
||||
|
||||
@@ -24,7 +24,10 @@ import java.util.Date;
|
||||
import org.apache.rocketmq.broker.BrokerController;
|
||||
import org.apache.rocketmq.broker.BrokerStartup;
|
||||
import org.apache.rocketmq.client.log.ClientLogger;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
@@ -34,11 +37,10 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcForwardService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.LocalGrpcService;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ProxyStartup {
|
||||
private static final Logger log = LoggerFactory.getLogger(ProxyStartup.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
private static final ProxyStartAndShutdown PROXY_START_AND_SHUTDOWN = new ProxyStartAndShutdown();
|
||||
|
||||
private static class ProxyStartAndShutdown extends AbstractStartAndShutdown {
|
||||
|
||||
@@ -27,14 +27,14 @@ import java.util.concurrent.ConcurrentMap;
|
||||
import java.util.function.Supplier;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.channel.GrpcClientChannel;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ChannelManager {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
private final ConcurrentMap<String /* clientId */, SimpleChannel> clientIdChannelMap = new ConcurrentHashMap<>();
|
||||
private final ConcurrentMap<String /* group */, Set<String>/* clientId */> groupClientIdMap = new ConcurrentHashMap<>();
|
||||
|
||||
|
||||
@@ -30,8 +30,8 @@ import io.netty.util.concurrent.GlobalEventExecutor;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.net.SocketAddress;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
|
||||
/**
|
||||
* SimpleChannel is used to handle writeAndFlush situation in processor
|
||||
@@ -39,7 +39,7 @@ import org.slf4j.LoggerFactory;
|
||||
* @see io.netty.channel.Channel#writeAndFlush
|
||||
*/
|
||||
public class SimpleChannel extends AbstractChannel {
|
||||
protected static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
protected final String remoteAddress;
|
||||
protected final String localAddress;
|
||||
|
||||
@@ -21,11 +21,12 @@ import com.alibaba.fastjson.JSON;
|
||||
import java.io.File;
|
||||
import java.nio.file.Files;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class Configuration {
|
||||
private final static Logger log = LoggerFactory.getLogger(Configuration.class);
|
||||
private final static Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
private final AtomicReference<ProxyConfig> proxyConfigReference = new AtomicReference<>();
|
||||
|
||||
public void init() throws Exception {
|
||||
|
||||
@@ -35,6 +35,7 @@ import org.apache.rocketmq.client.impl.CommunicationMode;
|
||||
import org.apache.rocketmq.client.impl.MQClientAPIImpl;
|
||||
import org.apache.rocketmq.client.impl.consumer.PullResultExt;
|
||||
import org.apache.rocketmq.client.producer.SendResult;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.message.Message;
|
||||
import org.apache.rocketmq.common.message.MessageBatch;
|
||||
import org.apache.rocketmq.common.message.MessageClientIDSetter;
|
||||
@@ -56,16 +57,16 @@ import org.apache.rocketmq.common.protocol.header.SearchOffsetResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.UpdateConsumerOffsetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.RPCHook;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingException;
|
||||
import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
import org.apache.rocketmq.remoting.netty.ResponseFuture;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class MQClientAPIExt extends MQClientAPIImpl {
|
||||
private static final Logger LOGGER = LoggerFactory.getLogger(MQClientAPIExt.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final ClientConfig clientConfig;
|
||||
|
||||
@@ -82,7 +83,7 @@ public class MQClientAPIExt extends MQClientAPIImpl {
|
||||
public boolean updateNameServerAddressList() {
|
||||
if (this.clientConfig.getNamesrvAddr() != null) {
|
||||
this.updateNameServerAddressList(this.clientConfig.getNamesrvAddr());
|
||||
LOGGER.info("user specified name server address: {}", this.clientConfig.getNamesrvAddr());
|
||||
log.info("user specified name server address: {}", this.clientConfig.getNamesrvAddr());
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
|
||||
+4
-3
@@ -19,13 +19,14 @@ package org.apache.rocketmq.proxy.connector.factory;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.RPCHook;
|
||||
import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public abstract class AbstractClientManager<T> {
|
||||
private static final Logger log = LoggerFactory.getLogger(AbstractClientManager.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
protected final ScheduledExecutorService scheduledExecutorService;
|
||||
protected Map<String, T> cacheTable = new ConcurrentHashMap<>();
|
||||
|
||||
@@ -24,19 +24,20 @@ import java.util.concurrent.ThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.client.exception.MQClientException;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.AbstractCacheLoader;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.config.ProxyConfig;
|
||||
import org.apache.rocketmq.proxy.connector.DefaultForwardClient;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class TopicRouteCache {
|
||||
private static final Logger log = LoggerFactory.getLogger(TopicRouteCache.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final LoadingCache<String /* topicName */, MessageQueueWrapper> topicCache;
|
||||
private final ThreadPoolExecutor cacheRefreshExecutor;
|
||||
|
||||
+4
-3
@@ -28,21 +28,22 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.common.ServiceThread;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.ProducerData;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.config.ProxyConfig;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardProducer;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteCache;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class TransactionHeartbeatRegisterService implements StartAndShutdown {
|
||||
private static final Logger log = LoggerFactory.getLogger(TransactionHeartbeatRegisterService.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private static final String TRANS_HEARTBEAT_CLIENT_ID = "rmq-proxy-producer-client";
|
||||
|
||||
|
||||
+4
-3
@@ -26,15 +26,16 @@ import java.util.Objects;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.client.producer.SendResult;
|
||||
import org.apache.rocketmq.common.UtilAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.message.MessageDecoder;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.message.MessageId;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.common.RemotingUtil;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class TransactionId {
|
||||
private static final Logger log = LoggerFactory.getLogger(TransactionId.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private SocketAddress brokerAddr;
|
||||
private String brokerTransactionId;
|
||||
|
||||
@@ -36,6 +36,8 @@ import org.apache.rocketmq.acl.AccessValidator;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
|
||||
import org.apache.rocketmq.common.utils.ServiceProvider;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.AuthenticationInterceptor;
|
||||
@@ -43,11 +45,9 @@ import org.apache.rocketmq.proxy.grpc.interceptor.ContextInterceptor;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.HeaderInterceptor;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.GrpcMessagingProcessor;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcForwardService;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class GrpcServer implements StartAndShutdown {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final io.grpc.Server server;
|
||||
private final ThreadPoolExecutor executor;
|
||||
|
||||
@@ -59,15 +59,15 @@ import io.grpc.stub.StreamObserver;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.CompletionException;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.ProxyException;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseWriter;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.service.GrpcForwardService;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseWriter;
|
||||
|
||||
public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServiceImplBase {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
private final GrpcForwardService grpcForwardService;
|
||||
|
||||
public GrpcMessagingProcessor(GrpcForwardService grpcForwardService) {
|
||||
|
||||
@@ -89,15 +89,15 @@ import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
import org.apache.rocketmq.common.sysflag.PullSysFlag;
|
||||
import org.apache.rocketmq.common.utils.BinaryUtil;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class GrpcConverter {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
public static String wrapResourceWithNamespace(Resource resource) {
|
||||
return NamespaceUtil.wrapNamespace(resource.getResourceNamespace(), resource.getName());
|
||||
|
||||
+3
-3
@@ -55,14 +55,14 @@ import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.V2Converter;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService clusterGrpcService;
|
||||
|
||||
|
||||
@@ -85,15 +85,15 @@ import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
import org.apache.rocketmq.common.utils.BinaryUtil;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.common.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class GrpcConverter {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
public static String wrapResourceWithNamespace(Resource resource) {
|
||||
return NamespaceUtil.wrapNamespace(resource.getResourceNamespace(), resource.getName());
|
||||
|
||||
@@ -20,11 +20,11 @@ package org.apache.rocketmq.proxy.grpc.v2.adapter;
|
||||
import io.grpc.stub.ServerCallStreamObserver;
|
||||
import io.grpc.stub.StreamObserver;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
|
||||
public class ResponseWriter {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
public static <T> void write(StreamObserver<T> observer, final T response) {
|
||||
if (observer instanceof ServerCallStreamObserver) {
|
||||
|
||||
+3
-3
@@ -36,16 +36,16 @@ import org.apache.rocketmq.common.message.MessageDecoder;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.protocol.header.ExtraInfoUtil;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.channel.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingSysResponseCode;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ReceiveMessageResponseHandler implements ResponseHandler<ReceiveMessageRequest, ReceiveMessageResponse> {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
private final String brokerName;
|
||||
private final boolean fifo;
|
||||
|
||||
|
||||
+3
-3
@@ -47,6 +47,8 @@ import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
@@ -60,11 +62,9 @@ import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ForwardClientService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ProducerService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.TransactionService;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(
|
||||
new ThreadFactoryImpl("ClusterGrpcServiceScheduledThread")
|
||||
|
||||
@@ -78,6 +78,8 @@ import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.UnregisterClientRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.channel.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
@@ -102,11 +104,9 @@ import org.apache.rocketmq.remoting.RemotingServer;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRemotingAbstract;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcForwardService {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final BrokerController brokerController;
|
||||
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(
|
||||
|
||||
+4
-3
@@ -20,16 +20,17 @@ import apache.rocketmq.v2.SendMessageRequest;
|
||||
import io.grpc.Context;
|
||||
import java.util.List;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.message.Message;
|
||||
import org.apache.rocketmq.common.message.MessageConst;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteCache;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class DefaultWriteQueueSelector implements WriteQueueSelector {
|
||||
private static final Logger log = LoggerFactory.getLogger(DefaultWriteQueueSelector.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
protected final TopicRouteCache topicRouteCache;
|
||||
|
||||
|
||||
+4
-3
@@ -38,8 +38,11 @@ import org.apache.rocketmq.broker.client.ProducerChangeListener;
|
||||
import org.apache.rocketmq.broker.client.ProducerGroupEvent;
|
||||
import org.apache.rocketmq.broker.client.ProducerManager;
|
||||
import org.apache.rocketmq.common.MQVersion;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
@@ -51,11 +54,9 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.ClientSettingsService;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ForwardClientService extends BaseService {
|
||||
private static final Logger log = LoggerFactory.getLogger(ForwardClientService.class);
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
|
||||
|
||||
private final ChannelManager channelManager;
|
||||
private final ConsumerManager consumerManager;
|
||||
|
||||
-12
@@ -164,16 +164,4 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
assertSame(ex, e.getCause());
|
||||
}
|
||||
}
|
||||
|
||||
// @Test
|
||||
// public void testForwardMessageToDeadLetterQueue() throws Exception {
|
||||
// ProducerService producerService = new ProducerService(this.connectorManager);
|
||||
//
|
||||
// when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr");
|
||||
// producerService.forwardMessageToDeadLetterQueue(Context.current(), ForwardMessageToDeadLetterQueueRequest.newBuilder()
|
||||
// .setMessageId("msgId")
|
||||
// .setReceiptHandle(createReceiptHandle().encode())
|
||||
// .set
|
||||
// .build());
|
||||
// }
|
||||
}
|
||||
Reference in New Issue
Block a user