diff --git a/agents/drivers/rocketmq/src/main/java/com/dbx/agent/rocketmq/RocketMqAgent.java b/agents/drivers/rocketmq/src/main/java/com/dbx/agent/rocketmq/RocketMqAgent.java index dbb1e1122..c39f120ae 100644 --- a/agents/drivers/rocketmq/src/main/java/com/dbx/agent/rocketmq/RocketMqAgent.java +++ b/agents/drivers/rocketmq/src/main/java/com/dbx/agent/rocketmq/RocketMqAgent.java @@ -43,6 +43,11 @@ import java.io.BufferedReader; import java.io.InputStreamReader; import java.nio.charset.StandardCharsets; import java.util.*; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; /** * RocketMQ admin agent for DBX. Communicates with the Rust bridge via JSON-RPC @@ -54,6 +59,8 @@ public final class RocketMqAgent { private static final Gson GSON = new GsonBuilder().serializeNulls().create(); private static final int DEFAULT_REQUEST_TIMEOUT_MS = 30_000; private static final int DEFAULT_LIST_LIMIT = 200; + private static final int CONSUMER_GROUP_ENRICH_CONCURRENCY = 8; + private static final long CONSUMER_GROUP_ENRICH_BUDGET_MS = 8_000; private static final String AUTO_CREATE_TOPIC_KEY = "TBW102"; /** RocketMQ 5.x topic attribute key for message type (NORMAL/DELAY/FIFO/TRANSACTION). */ private static final String TOPIC_MESSAGE_TYPE_ATTRIBUTE = "message.type"; @@ -946,6 +953,25 @@ public final class RocketMqAgent { return null; } + private static Map collectConsumerGroupConfigs( + DefaultMQAdminExt admin, JsonObject conn) throws Exception { + Map configs = new TreeMap<>(); + for (String brokerAddr : resolveMasterBrokerAddrs(admin, conn)) { + try { + SubscriptionGroupWrapper wrapper = admin.getAllSubscriptionGroup(brokerAddr, DEFAULT_REQUEST_TIMEOUT_MS); + if (wrapper == null || wrapper.getSubscriptionGroupTable() == null) { + continue; + } + for (Map.Entry entry : wrapper.getSubscriptionGroupTable().entrySet()) { + configs.putIfAbsent(entry.getKey(), entry.getValue()); + } + } catch (Exception ignored) { + // Try next broker. + } + } + return configs; + } + /** * DefaultMQAdminExt.examineConsumerConnectionInfo(group) picks a broker from route using * NameServer-registered addresses. Remap to the client-reachable host before querying. @@ -1110,11 +1136,12 @@ public final class RocketMqAgent { limit = DEFAULT_LIST_LIMIT; } + Map configs = collectConsumerGroupConfigs(admin, conn); Set groups = new TreeSet<>(); if (!topicFilter.isBlank()) { groups.addAll(queryTopicConsumeByWhoRemapped(admin, conn, topicFilter)); } else { - groups.addAll(collectAllConsumerGroups(admin, conn)); + groups.addAll(configs.keySet()); } List> rows = new ArrayList<>(); @@ -1135,13 +1162,11 @@ public final class RocketMqAgent { List> page = paginate(rows, offset, limit); for (Map row : page) { String groupId = String.valueOf(row.get("groupId")); - SubscriptionGroupConfig config = findSubscriptionGroupConfig(admin, conn, groupId); + SubscriptionGroupConfig config = configs.get(groupId); row.put("groupType", classifyConsumerGroupType(groupId, config)); } if (boolOrDefault(params, "enrich", false)) { - for (Map row : page) { - enrichConsumerGroupRow(admin, conn, String.valueOf(row.get("groupId")), row); - } + enrichConsumerGroupRows(admin, conn, page); } Map result = new LinkedHashMap<>(); @@ -1152,18 +1177,72 @@ public final class RocketMqAgent { return result; } - private static void enrichConsumerGroupRow( - DefaultMQAdminExt admin, JsonObject conn, String groupId, Map row) { + @FunctionalInterface + interface ConsumerConnectionLookup { + ConsumerConnection load(String groupId) throws Exception; + } + + private static void enrichConsumerGroupRows( + DefaultMQAdminExt admin, JsonObject conn, List> page) { + enrichConsumerGroupRows( + page, + groupId -> examineConsumerConnectionInfoRemapped(admin, conn, groupId), + CONSUMER_GROUP_ENRICH_CONCURRENCY, + CONSUMER_GROUP_ENRICH_BUDGET_MS + ); + } + + static void enrichConsumerGroupRows( + List> page, + ConsumerConnectionLookup lookup, + int concurrency, + long budgetMs) { + if (page.isEmpty()) { + return; + } + int workers = Math.max(1, Math.min(concurrency, page.size())); + ExecutorService executor = Executors.newFixedThreadPool(workers, runnable -> { + Thread thread = new Thread(runnable, "dbx-rocketmq-consumer-enrich"); + thread.setDaemon(true); + return thread; + }); + List>> tasks = new ArrayList<>(page.size()); + for (Map row : page) { + tasks.add(() -> consumerGroupEnrichment(String.valueOf(row.get("groupId")), lookup)); + } + List>> results = Collections.emptyList(); try { - SubscriptionGroupConfig config = findSubscriptionGroupConfig(admin, conn, groupId); - row.put("groupType", classifyConsumerGroupType(groupId, config)); - ConsumerConnection connection = examineConsumerConnectionInfoRemapped(admin, conn, groupId); - row.put("consumeType", connection.getConsumeType() != null ? connection.getConsumeType().name() : "UNKNOWN"); + results = executor.invokeAll(tasks, Math.max(1, budgetMs), TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + executor.shutdownNow(); + } + for (int index = 0; index < page.size(); index++) { + Map row = page.get(index); + if (index < results.size() && !results.get(index).isCancelled()) { + try { + row.putAll(results.get(index).get()); + } catch (Exception ignored) { + // Keep the base group row when enrichment fails. + } + } + row.putIfAbsent("memberCount", 0); + row.putIfAbsent("topics", Collections.emptyList()); + } + } + + private static Map consumerGroupEnrichment( + String groupId, ConsumerConnectionLookup lookup) { + Map enrichment = new LinkedHashMap<>(); + try { + ConsumerConnection connection = lookup.load(groupId); + enrichment.put("consumeType", connection.getConsumeType() != null ? connection.getConsumeType().name() : "UNKNOWN"); if (connection.getMessageModel() != null) { - row.put("messageModel", connection.getMessageModel().name()); + enrichment.put("messageModel", connection.getMessageModel().name()); } int memberCount = connection.getConnectionSet() == null ? 0 : connection.getConnectionSet().size(); - row.put("memberCount", memberCount); + enrichment.put("memberCount", memberCount); List topics = new ArrayList<>(); if (connection.getSubscriptionTable() != null) { for (SubscriptionData sub : connection.getSubscriptionTable().values()) { @@ -1172,11 +1251,12 @@ public final class RocketMqAgent { } } } - row.put("topics", topics); + enrichment.put("topics", topics); } catch (Exception ignored) { - row.putIfAbsent("memberCount", 0); - row.putIfAbsent("topics", Collections.emptyList()); + enrichment.put("memberCount", 0); + enrichment.put("topics", Collections.emptyList()); } + return enrichment; } private static Object describeConsumerGroup(JsonObject params) throws Exception { @@ -1747,21 +1827,6 @@ public final class RocketMqAgent { return admin.fetchAllTopicList(); } - private static Set collectAllConsumerGroups(DefaultMQAdminExt admin, JsonObject conn) throws Exception { - Set groups = new TreeSet<>(); - for (String brokerAddr : resolveMasterBrokerAddrs(admin, conn)) { - try { - SubscriptionGroupWrapper wrapper = admin.getAllSubscriptionGroup(brokerAddr, DEFAULT_REQUEST_TIMEOUT_MS); - if (wrapper != null && wrapper.getSubscriptionGroupTable() != null) { - groups.addAll(wrapper.getSubscriptionGroupTable().keySet()); - } - } catch (Exception ignored) { - // Try next broker when Docker/internal broker addresses are unreachable. - } - } - return groups; - } - private static TopicConfig loadTopicConfig(DefaultMQAdminExt admin, String brokerAddr, String topic) throws Exception { TopicConfig config = admin.examineTopicConfig(brokerAddr, topic); diff --git a/agents/drivers/rocketmq/src/test/java/com/dbx/agent/rocketmq/RocketMqAgentTest.java b/agents/drivers/rocketmq/src/test/java/com/dbx/agent/rocketmq/RocketMqAgentTest.java index 00192040f..e0531cb32 100644 --- a/agents/drivers/rocketmq/src/test/java/com/dbx/agent/rocketmq/RocketMqAgentTest.java +++ b/agents/drivers/rocketmq/src/test/java/com/dbx/agent/rocketmq/RocketMqAgentTest.java @@ -14,6 +14,7 @@ import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable; import org.apache.rocketmq.remoting.protocol.admin.TopicOffset; import org.apache.rocketmq.remoting.protocol.body.ProducerInfo; import org.apache.rocketmq.remoting.protocol.body.ProducerTableInfo; +import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection; import org.apache.rocketmq.common.message.MessageQueue; import com.google.gson.JsonObject; import com.google.gson.JsonParser; @@ -23,6 +24,10 @@ import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.IntStream; import org.junit.jupiter.api.Test; @@ -223,6 +228,70 @@ class RocketMqAgentTest { assertEquals("NORMAL", RocketMqAgent.classifyConsumerGroupType("MyGroup", normal)); } + @Test + void enrichConsumerGroupRowsRunsIndependentLookupsConcurrently() throws Exception { + List> rows = IntStream.range(0, 4) + .mapToObj(index -> { + Map row = new LinkedHashMap<>(); + row.put("groupId", "group-" + index); + return row; + }) + .toList(); + CountDownLatch started = new CountDownLatch(rows.size()); + CountDownLatch release = new CountDownLatch(1); + AtomicInteger active = new AtomicInteger(); + AtomicInteger maxActive = new AtomicInteger(); + + CompletableFuture enrichment = CompletableFuture.runAsync(() -> + RocketMqAgent.enrichConsumerGroupRows(rows, groupId -> { + int current = active.incrementAndGet(); + maxActive.accumulateAndGet(current, Math::max); + started.countDown(); + try { + release.await(1, TimeUnit.SECONDS); + return new ConsumerConnection(); + } finally { + active.decrementAndGet(); + } + }, 4, 1_000) + ); + + assertTrue(started.await(1, TimeUnit.SECONDS)); + release.countDown(); + enrichment.get(2, TimeUnit.SECONDS); + + assertTrue(maxActive.get() > 1); + for (Map row : rows) { + assertEquals(0, row.get("memberCount")); + assertEquals(List.of(), row.get("topics")); + } + } + + @Test + void enrichConsumerGroupRowsReturnsDefaultsWhenBudgetExpires() { + List> rows = IntStream.range(0, 4) + .mapToObj(index -> { + Map row = new LinkedHashMap<>(); + row.put("groupId", "slow-group-" + index); + return row; + }) + .toList(); + CountDownLatch blocked = new CountDownLatch(1); + long startedAt = System.nanoTime(); + + RocketMqAgent.enrichConsumerGroupRows(rows, groupId -> { + blocked.await(); + return new ConsumerConnection(); + }, 2, 50); + + long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startedAt); + assertTrue(elapsedMs < 1_000, "enrichment exceeded its response budget: " + elapsedMs + "ms"); + for (Map row : rows) { + assertEquals(0, row.get("memberCount")); + assertEquals(List.of(), row.get("topics")); + } + } + @Test void isEmptyQueryMessageResultDetectsRocketMqCode208() { MQClientException empty = new MQClientException(208, "query message by key finished, but no message");