This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new fae6c5a9 fix(message): keep newest topic query messages across all
queues (#2043)
fae6c5a9 is described below
commit fae6c5a9cc13c986c470e0c4c14bb5339b9c0f74
Author: xdz997 <[email protected]>
AuthorDate: Fri Aug 14 10:39:22 2026 +0800
fix(message): keep newest topic query messages across all queues (#2043)
Replace the early-stop collection with a priority queue so messages from
every queue compete for the limit by store timestamp, instead of stopping
once earlier queues fill the page (follow-up to the time-sort fix).
---
.../provider/apache/RocketMQMessageProvider.java | 39 ++++++---
.../apache/RocketMQMessageProviderTest.java | 94 ++++++++++++++++++----
2 files changed, 106 insertions(+), 27 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index 8d315b20..1432e823 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -56,6 +56,7 @@ import java.util.Base64;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.PriorityQueue;
import java.util.Set;
/**
@@ -81,6 +82,9 @@ public class RocketMQMessageProvider implements
MessageProvider {
private static final long ONE_HOUR_MILLIS = 3600_000L;
private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
private static final int MAX_CONSECUTIVE_OFFSET_ILLEGAL = 3;
+ private static final Comparator<MessageRecordVO> TOPIC_QUERY_ORDER =
Comparator
+ .comparingLong(MessageRecordVO::getStoreTime)
+ .thenComparing(MessageRecordVO::getMsgId,
Comparator.nullsFirst(String::compareTo));
private final RuntimeAdminClientResolver runtimeAdminClientResolver;
private final QueryHistoryService queryHistoryService;
@@ -191,22 +195,19 @@ public class RocketMQMessageProvider implements
MessageProvider {
*/
private List<MessageRecordVO> queryByTopic(String endpoint, String topic,
String tag, long begin, long end, int limit) {
DefaultMQPullConsumer consumer = newPullConsumer("studio-msg-query",
endpoint);
- List<MessageRecordVO> result = new ArrayList<>();
+ int resultLimit = Math.min(limit, TOPIC_QUERY_HARD_CAP);
+ PriorityQueue<MessageRecordVO> newestMessages = new
PriorityQueue<>(TOPIC_QUERY_ORDER);
try {
consumer.start();
Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(topic);
if (queues == null || queues.isEmpty()) {
- return result;
+ return Collections.emptyList();
}
- outer:
for (MessageQueue queue : queues) {
long minOffset = consumer.searchOffset(queue, begin);
long maxOffset = consumer.searchOffset(queue, end);
int consecutiveIllegalOffsets = 0;
for (long offset = minOffset; offset <= maxOffset; ) {
- if (result.size() >= Math.min(limit,
TOPIC_QUERY_HARD_CAP)) {
- break outer;
- }
PullResult pullResult = consumer.pull(queue, "*", offset,
32);
if (pullResult == null) {
log.warn("Stop topic query for {} because queue {}
returned no pull result", topic, queue);
@@ -247,10 +248,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
if (!matchesTag(messageExt, tag)) {
continue;
}
- result.add(toRecordVO(messageExt));
- if (result.size() >= Math.min(limit,
TOPIC_QUERY_HARD_CAP)) {
- break outer;
- }
+ addTopicQueryCandidate(newestMessages,
toRecordVO(messageExt), resultLimit);
}
}
}
@@ -260,8 +258,25 @@ public class RocketMQMessageProvider implements
MessageProvider {
} finally {
consumer.shutdown();
}
-
result.sort(Comparator.comparingLong(MessageRecordVO::getStoreTime).reversed());
- return result;
+ return newestMessages.stream()
+ .sorted(TOPIC_QUERY_ORDER.reversed())
+ .toList();
+ }
+
+ private void addTopicQueryCandidate(PriorityQueue<MessageRecordVO>
newestMessages,
+ MessageRecordVO candidate, int
resultLimit) {
+ if (resultLimit <= 0) {
+ return;
+ }
+ if (newestMessages.size() < resultLimit) {
+ newestMessages.offer(candidate);
+ return;
+ }
+ MessageRecordVO oldestKept = newestMessages.peek();
+ if (oldestKept != null && TOPIC_QUERY_ORDER.compare(candidate,
oldestKept) > 0) {
+ newestMessages.poll();
+ newestMessages.offer(candidate);
+ }
}
@Override
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 87bbed05..033241dc 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -45,9 +45,11 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
+import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.TimeUnit;
+import java.util.stream.IntStream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -193,24 +195,21 @@ class RocketMQMessageProviderTest {
}
@Test
- void queryByTopicReturnsNewestMessagesFirst() throws Exception {
- MessageQueue queue = new MessageQueue("TopicA", "broker-a", 0);
- MessageExt older = new MessageExt();
- older.setMsgId("older");
- older.setTopic("TopicA");
- older.setStoreTimestamp(150L);
- MessageExt newer = new MessageExt();
- newer.setMsgId("newer");
- newer.setTopic("TopicA");
- newer.setStoreTimestamp(250L);
- PullResult pullResult = new PullResult(PullStatus.FOUND, 11L, 10L, 11L,
- List.of(older, newer));
+ void queryByTopicSortsMessagesAcrossQueuesNewestFirst() throws Exception {
+ MessageQueue olderQueue = new MessageQueue("TopicA", "broker-a", 0);
+ MessageQueue newerQueue = new MessageQueue("TopicA", "broker-a", 1);
+ PullResult olderPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
+ List.of(topicMessage("older", 150L)));
+ PullResult newerPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
+ List.of(topicMessage("newer", 250L)));
try (MockedConstruction<DefaultMQPullConsumer> ignored =
mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
- when(consumer.searchOffset(eq(queue),
anyLong())).thenReturn(10L);
- when(consumer.pull(queue, "*", 10L,
32)).thenReturn(pullResult);
+ when(consumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenReturn(new
LinkedHashSet<>(List.of(olderQueue, newerQueue)));
+ when(consumer.searchOffset(any(MessageQueue.class),
anyLong())).thenReturn(10L);
+ when(consumer.pull(olderQueue, "*", 10L,
32)).thenReturn(olderPullResult);
+ when(consumer.pull(newerQueue, "*", 10L,
32)).thenReturn(newerPullResult);
doNothing().when(consumer).shutdown();
})) {
List<MessageRecordVO> messages = provider.queryMessages(
@@ -256,6 +255,62 @@ class RocketMQMessageProviderTest {
100L, 200L, 1);
}
+ @Test
+ void queryByTopicKeepsNewestMessagesWhenEarlierQueuesFillTheDefaultLimit()
throws Exception {
+ MessageQueue olderQueue = new MessageQueue("TopicA", "broker-a", 0);
+ MessageQueue newerQueue = new MessageQueue("TopicA", "broker-a", 1);
+ PullResult olderPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
+ IntStream.range(0, 200)
+ .mapToObj(index -> topicMessage("older-" + index,
150L))
+ .toList());
+ PullResult newerPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
+ List.of(topicMessage("newer", 250L)));
+ try (MockedConstruction<DefaultMQPullConsumer> ignored =
+ mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
+ doNothing().when(consumer).start();
+ when(consumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenReturn(new
LinkedHashSet<>(List.of(olderQueue, newerQueue)));
+ when(consumer.searchOffset(any(MessageQueue.class),
anyLong())).thenReturn(10L);
+ when(consumer.pull(olderQueue, "*", 10L,
32)).thenReturn(olderPullResult);
+ when(consumer.pull(newerQueue, "*", 10L,
32)).thenReturn(newerPullResult);
+ doNothing().when(consumer).shutdown();
+ })) {
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+ assertThat(messages).hasSize(200);
+ assertThat(messages.get(0).getMsgId()).isEqualTo("newer");
+ assertThat(messages).extracting(MessageRecordVO::getMsgId)
+ .contains("newer");
+ }
+ }
+
+ @Test
+ void queryByTopicUsesMessageIdAsStableTieBreakerForEqualTimestamps()
throws Exception {
+ MessageQueue firstQueue = new MessageQueue("TopicA", "broker-a", 0);
+ MessageQueue secondQueue = new MessageQueue("TopicA", "broker-a", 1);
+ PullResult firstPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
+ List.of(topicMessage("message-a", 250L)));
+ PullResult secondPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
+ List.of(topicMessage("message-b", 250L)));
+ try (MockedConstruction<DefaultMQPullConsumer> ignored =
+ mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
+ doNothing().when(consumer).start();
+ when(consumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenReturn(new
LinkedHashSet<>(List.of(firstQueue, secondQueue)));
+ when(consumer.searchOffset(any(MessageQueue.class),
anyLong())).thenReturn(10L);
+ when(consumer.pull(firstQueue, "*", 10L,
32)).thenReturn(firstPullResult);
+ when(consumer.pull(secondQueue, "*", 10L,
32)).thenReturn(secondPullResult);
+ doNothing().when(consumer).shutdown();
+ })) {
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+ assertThat(messages).extracting(MessageRecordVO::getMsgId)
+ .containsExactly("message-b", "message-a");
+ }
+ }
+
@Test
void toRecordVOBoundsMessageBodyAndProperties() {
MessageExt message = new MessageExt();
@@ -436,4 +491,13 @@ class RocketMQMessageProviderTest {
return String.join(String.valueOf(TraceConstants.FIELD_SPLITOR),
contexts)
+ TraceConstants.FIELD_SPLITOR;
}
+
+ private static MessageExt topicMessage(String msgId, long storeTimestamp) {
+ MessageExt message = new MessageExt();
+ message.setMsgId(msgId);
+ message.setTopic("TopicA");
+ message.setStoreTimestamp(storeTimestamp);
+ message.setBody(("body-" + msgId).getBytes(StandardCharsets.UTF_8));
+ return message;
+ }
}