mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Improve readability.
This commit is contained in:
@@ -503,10 +503,8 @@ public class Converter {
|
||||
public static Message buildMessage(MessageExt messageExt) {
|
||||
Map<String, String> userAttributes = buildUserAttributes(messageExt);
|
||||
SystemAttribute systemAttributes = buildSystemAttributes(messageExt);
|
||||
Resource topic = Resource.newBuilder()
|
||||
.setResourceNamespace(NamespaceUtil.getNamespaceFromResource(messageExt.getTopic()))
|
||||
.setName(NamespaceUtil.withoutNamespace(messageExt.getTopic()))
|
||||
.build();
|
||||
Resource topic = buildResource(messageExt.getTopic());
|
||||
|
||||
return Message.newBuilder()
|
||||
.setTopic(topic)
|
||||
.putAllUserAttribute(userAttributes)
|
||||
@@ -641,12 +639,7 @@ public class Converter {
|
||||
// publisher_group
|
||||
String producerGroup = messageExt.getProperty(MessageConst.PROPERTY_PRODUCER_GROUP);
|
||||
if (producerGroup != null) {
|
||||
String namespaceId = NamespaceUtil.getNamespaceFromResource(producerGroup);
|
||||
String group = NamespaceUtil.withoutNamespace(producerGroup);
|
||||
systemAttributeBuilder.setProducerGroup(Resource.newBuilder()
|
||||
.setResourceNamespace(namespaceId)
|
||||
.setName(group)
|
||||
.build());
|
||||
systemAttributeBuilder.setProducerGroup(buildResource(producerGroup));
|
||||
}
|
||||
|
||||
// trace context
|
||||
@@ -690,6 +683,13 @@ public class Converter {
|
||||
return consumeMessageDirectlyResult;
|
||||
}
|
||||
|
||||
public static Resource buildResource(String resourceNameWithNamespace) {
|
||||
return Resource.newBuilder()
|
||||
.setResourceNamespace(NamespaceUtil.getNamespaceFromResource(resourceNameWithNamespace))
|
||||
.setName(NamespaceUtil.withoutNamespace(resourceNameWithNamespace))
|
||||
.build();
|
||||
}
|
||||
|
||||
public static UnregisterClientRequestHeader buildUnregisterClientRequestHeader(NotifyClientTerminationRequest request) {
|
||||
UnregisterClientRequestHeader header = new UnregisterClientRequestHeader();
|
||||
header.setClientID(request.getClientId());
|
||||
|
||||
+11
-6
@@ -24,6 +24,9 @@ import apache.rocketmq.v1.PollCommandRequest;
|
||||
import apache.rocketmq.v1.PollCommandResponse;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import io.grpc.Context;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.broker.client.ClientChannelInfo;
|
||||
import org.apache.rocketmq.broker.client.ConsumerManager;
|
||||
import org.apache.rocketmq.broker.client.ProducerManager;
|
||||
@@ -32,16 +35,12 @@ import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
public class ClientService extends BaseService {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ClientService.class);
|
||||
@@ -52,7 +51,12 @@ public class ClientService extends BaseService {
|
||||
private final ProducerManager producerManager;
|
||||
private final PollCommandResponseManager pollCommandResponseManager;
|
||||
|
||||
public ClientService(ConnectorManager connectorManager, ScheduledExecutorService scheduledExecutorService, ChannelManager channelManager, PollCommandResponseManager pollCommandResponseManager) {
|
||||
public ClientService(
|
||||
ConnectorManager connectorManager,
|
||||
ScheduledExecutorService scheduledExecutorService,
|
||||
ChannelManager channelManager,
|
||||
PollCommandResponseManager pollCommandResponseManager
|
||||
) {
|
||||
super(connectorManager);
|
||||
scheduledExecutorService.scheduleWithFixedDelay(this::scanNotActiveChannel, 1000 * 10, 1000 * 10, TimeUnit.MILLISECONDS);
|
||||
this.channelManager = channelManager;
|
||||
@@ -70,6 +74,7 @@ public class ClientService extends BaseService {
|
||||
if (request.hasProducerData()) {
|
||||
String producerGroup = Converter.getResourceNameWithNamespace(request.getProducerData().getGroup());
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, producerGroup, clientId, pollCommandResponseManager);
|
||||
//TODO: Use the faked MQ Version ?
|
||||
ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal());
|
||||
producerManager.registerProducer(producerGroup, clientChannelInfo);
|
||||
}
|
||||
|
||||
+2
-1
@@ -34,7 +34,8 @@ public class DefaultAssignmentQueueSelector implements AssignmentQueueSelector {
|
||||
|
||||
@Override
|
||||
public List<SelectableMessageQueue> getAssignment(Context ctx, QueryAssignmentRequest request) throws Exception {
|
||||
MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic()));
|
||||
String topicName = Converter.getResourceNameWithNamespace(request.getTopic());
|
||||
MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(topicName);
|
||||
return messageQueueWrapper.getReadSelector().getBrokerActingQueues();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user