diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java index ea6f253f8e..f56627a257 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java @@ -19,6 +19,7 @@ package org.apache.rocketmq.proxy.processor; import io.netty.channel.Channel; import java.util.Set; import org.apache.rocketmq.broker.client.ClientChannelInfo; +import org.apache.rocketmq.broker.client.ConsumerGroupInfo; import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.broker.client.ProducerChangeListener; import org.apache.rocketmq.common.consumer.ConsumeFromWhere; @@ -101,4 +102,8 @@ public class ClientProcessor extends AbstractProcessor { public void registerConsumerIdsChangeListener(ConsumerIdsChangeListener listener) { this.serviceManager.getConsumerManager().appendConsumerIdsChangeListener(listener); } + + public ConsumerGroupInfo getConsumerGroupInfo(String consumerGroup) { + return this.serviceManager.getConsumerManager().getConsumerGroupInfo(consumerGroup); + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java index 628f1af73d..5e5ce2e39b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java @@ -24,6 +24,7 @@ import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.client.ClientChannelInfo; +import org.apache.rocketmq.broker.client.ConsumerGroupInfo; import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.broker.client.ProducerChangeListener; import org.apache.rocketmq.client.consumer.AckResult; @@ -217,6 +218,11 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen this.clientProcessor.registerConsumerIdsChangeListener(consumerIdsChangeListener); } + @Override + public ConsumerGroupInfo getConsumerGroupInfo(String consumerGroup) { + return this.clientProcessor.getConsumerGroupInfo(consumerGroup); + } + @Override public void addTransactionSubscription(ProxyContext ctx, String producerGroup, String topic) { this.transactionProcessor.addTransactionSubscription(ctx, producerGroup, topic); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java index c01e9c882a..2c957bbf5e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java @@ -22,6 +22,7 @@ import java.util.List; import java.util.Set; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.broker.client.ClientChannelInfo; +import org.apache.rocketmq.broker.client.ConsumerGroupInfo; import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.broker.client.ProducerChangeListener; import org.apache.rocketmq.client.consumer.AckResult; @@ -222,6 +223,8 @@ public interface MessagingProcessor extends StartAndShutdown { ConsumerIdsChangeListener consumerIdsChangeListener ); + ConsumerGroupInfo getConsumerGroupInfo(String consumerGroup); + void addTransactionSubscription( ProxyContext ctx, String producerGroup,