mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Add getConsumerGroupInfo
This commit is contained in:
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user