From 65fbd737cc467d7677c076bccb267b30ffeacda4 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 00:04:06 -0700 Subject: [PATCH 1/2] [ISSUE #10778] Skip empty priority queue groups --- .../service/route/MessageQueuePenalizer.java | 14 +++++--- .../route/MessageQueuePenalizerTest.java | 34 +++++++++++++++++++ 2 files changed, 44 insertions(+), 4 deletions(-) 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..b26a51e2041 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,9 @@ static Pair selectLeastPenaltyWithPriority( bestQueue = queueAndPenalty.getLeft(); } } - return Pair.of(bestQueue, bestPenalty); + if (bestQueue == null) { + 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 */ From 799a8aeb948ac669ed02e4d98a716b6a00858962 Mon Sep 17 00:00:00 2001 From: liuhy Date: Fri, 14 Aug 2026 06:25:46 -0700 Subject: [PATCH 2/2] docs(proxy): explain empty priority group result --- .../rocketmq/proxy/service/route/MessageQueuePenalizer.java | 1 + 1 file changed, 1 insertion(+) 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 b26a51e2041..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 @@ -133,6 +133,7 @@ static Pair selectLeastPenaltyWithPriority( } } if (bestQueue == null) { + // All priority groups were empty or absent. return null; } return Pair.of(bestQueue, bestPenalty);