mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] refactor SelectableMessageQueue and do the merging work.
This commit is contained in:
@@ -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 {
|
||||
|
||||
+3
-21
@@ -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<AddressableMessageQueue> 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<AddressableMessageQueue> getQueues() {
|
||||
return queues;
|
||||
}
|
||||
|
||||
public List<AddressableMessageQueue> getBrokers() {
|
||||
return brokers;
|
||||
public List<AddressableMessageQueue> getBrokerActingQueues() {
|
||||
return brokerActingQueues;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+1
-1
@@ -29,6 +29,6 @@ public class ConsumerService extends BaseService {
|
||||
}
|
||||
|
||||
public CompletableFuture<ReceiveMessageResponse> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -70,7 +70,7 @@ public class RouteService extends BaseService {
|
||||
public List<AddressableMessageQueue> getAssignment(QueryAssignmentRequest request) throws Exception {
|
||||
MessageQueueWrapper messageQueueWrapper = clientManager.getTopicRouteCache()
|
||||
.getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic()));
|
||||
return messageQueueWrapper.getRead().getBrokers();
|
||||
return messageQueueWrapper.getRead().getBrokerActingQueues();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user