diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizer.java index d53056971dc..2301c28e919 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizer.java @@ -120,8 +120,11 @@ static Pair selectLeastPenaltyWithPriority( int bestPenalty = Integer.MAX_VALUE; for (List queues : queuesWithPriority) { Pair queueAndPenalty = selectLeastPenalty(queues, penalizers, startIndex); - int penalty = queueAndPenalty.getRight(); - if (queueAndPenalty.getRight() <= 0) { + if (queueAndPenalty == null) { + continue; + } + int penalty = queueAndPenalty.getRight(); + if (penalty <= 0) { return queueAndPenalty; } if (penalty < bestPenalty) { @@ -129,6 +132,10 @@ static Pair selectLeastPenaltyWithPriority( bestQueue = queueAndPenalty.getLeft(); } } - return Pair.of(bestQueue, bestPenalty); + if (bestQueue == null) { + // All priority groups were empty or absent. + return null; + } + return Pair.of(bestQueue, bestPenalty); } -} \ No newline at end of file +} diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizerTest.java index f31d973cce5..1594b3f5aa7 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueuePenalizerTest.java @@ -282,6 +282,40 @@ public void testSelectLeastPenaltyWithPriority_EmptyQueues() { assertNull(result); } + /** + * Test selectLeastPenaltyWithPriority with only empty priority groups should return null + */ + @Test + public void testSelectLeastPenaltyWithPriority_AllEmptyPriorityGroups() { + List> penalizers = Collections.singletonList(mq -> 10); + AtomicInteger startIndex = new AtomicInteger(0); + Pair result = MessageQueuePenalizer.selectLeastPenaltyWithPriority( + Arrays.asList(Collections.emptyList(), Collections.emptyList()), penalizers, startIndex); + assertNull(result); + } + + /** + * Test selectLeastPenaltyWithPriority skips empty priority groups and selects from non-empty groups + */ + @Test + public void testSelectLeastPenaltyWithPriority_SkipEmptyPriorityGroup() { + MessageQueue mq0 = new MessageQueue("topic", "broker", 0); + MessageQueue mq1 = new MessageQueue("topic", "broker", 1); + List queues = Arrays.asList(mq0, mq1); + + List> penalizers = Collections.singletonList( + mq -> mq.getQueueId() == 0 ? 20 : 10 + ); + + AtomicInteger startIndex = new AtomicInteger(0); + Pair result = MessageQueuePenalizer.selectLeastPenaltyWithPriority( + Arrays.asList(Collections.emptyList(), queues), penalizers, startIndex); + + assertNotNull(result); + assertEquals(mq1, result.getLeft()); + assertEquals(10, result.getRight().intValue()); + } + /** * Test selectLeastPenaltyWithPriority with single priority group delegates to selectLeastPenalty */