From be512d57470dbf61d84b8339bb831802e99ea978 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Tue, 15 Mar 2022 15:32:03 +0800 Subject: [PATCH] [ISSUE #3949] refactor SelectableMessageQueue and do the merging work. --- .../proxy/client/TopicRouteCache.java | 2 +- .../client/route/SelectableMessageQueue.java | 24 +++---------------- .../grpc/service/cluster/ConsumerService.java | 2 +- .../grpc/service/cluster/RouteService.java | 2 +- 4 files changed, 6 insertions(+), 24 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java index c7f911f685..dec5496668 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java @@ -71,7 +71,7 @@ public class TopicRouteCache { if (last == null) { return getMessageQueue(topic).getWrite().selectOne(false); } - return getMessageQueue(topic).getWrite().selectNextOne(last); + return getMessageQueue(topic).getWrite().selectNextQueue(last); } public AddressableMessageQueue selectOneWriteQueue(String topic, String brokerName, int queueId) throws Exception { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java index 728071e4f3..88f7b3cee8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java @@ -119,7 +119,7 @@ public class SelectableMessageQueue { AddressableMessageQueue mq = new AddressableMessageQueue( new MessageQueue(topicRoute.getTopicName(), qd.getBrokerName(), i), brokerAddr); - queueSet.add(mq); + queueSet.add(mq); } } } @@ -128,8 +128,6 @@ public class SelectableMessageQueue { return queueSet.stream().sorted().collect(Collectors.toList()); } - - // 这里应该不是线程安全的 private void buildBrokerActingQueues(String topic, List normalQueues) { for (AddressableMessageQueue mq : normalQueues) { AddressableMessageQueue brokerActingQueue = new AddressableMessageQueue( @@ -194,28 +192,12 @@ public class SelectableMessageQueue { return newOne; } - // should use selectNextQueue -// public final AddressableMessageQueue selectNextBrokerActingQueue(AddressableMessageQueue last) { -// boolean onlyBroker = last.getQueueId() < 0; -// AddressableMessageQueue newOne = last; -// int count = onlyBroker ? brokerActingQueues.size() : queues.size(); -// -// for (int i = 0; i < count; i++) { -// newOne = selectOne(onlyBroker); -// if (!newOne.getBrokerName().equals(last.getBrokerName())) { -// break; -// } -// } -// -// return newOne; -// } - public List getQueues() { return queues; } - public List getBrokers() { - return brokers; + public List getBrokerActingQueues() { + return brokerActingQueues; } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index 8894c16fb7..7aa6887443 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -29,6 +29,6 @@ public class ConsumerService extends BaseService { } public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { - + return null; } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index 4bcea12cc0..6bb10534bc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -70,7 +70,7 @@ public class RouteService extends BaseService { public List getAssignment(QueryAssignmentRequest request) throws Exception { MessageQueueWrapper messageQueueWrapper = clientManager.getTopicRouteCache() .getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic())); - return messageQueueWrapper.getRead().getBrokers(); + return messageQueueWrapper.getRead().getBrokerActingQueues(); } }