如何通过Kafka消费者ID查询其关联的topic列表
Kafka消费者关联Topic查询实现方案
核心逻辑基于Kafka官方提供的AdminClient接口实现,无需额外依赖第三方组件,即可直接从Kafka集群获取消费者订阅的Topic列表。
前置说明
通常你提到的「消费者ID」分为两种场景,对应不同的查询逻辑:
- 若为消费者组ID:直接查询整个消费组订阅的所有Topic即可
- 若为消费组内单个消费者实例ID:需要先查询消费组下所有实例的订阅信息,匹配到目标实例后再提取关联Topic
Java 实现方案
1. 引入依赖(Maven示例)
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>对应你的Kafka集群版本</version> </dependency>
2. 核心查询代码
import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.ListConsumerGroupOffsetsResult; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import java.util.Properties; import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; public class ConsumerTopicQuery { // 替换为你的Kafka集群地址 private static final String BOOTSTRAP_SERVERS = "kafka1:9092,kafka2:9092"; public static Set<String> getTopicsByConsumerGroup(String groupId) throws ExecutionException, InterruptedException { Properties props = new Properties(); props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); // 开了ACL的集群需要额外配置用户名密码等参数 // props.put("sasl.jaas.config", "xxx"); // props.put("security.protocol", "SASL_PLAINTEXT"); try (Admin admin = Admin.create(props)) { ListConsumerGroupOffsetsResult offsetsResult = admin.listConsumerGroupOffsets(groupId); Set<TopicPartition> partitions = offsetsResult.partitionsToOffsetAndMetadata().get().keySet(); // 对Topic去重返回 return partitions.stream().map(TopicPartition::topic).collect(Collectors.toSet()); } } // 如果是查询单个消费者实例关联的Topic public static Set<String> getTopicsByConsumerInstanceId(String groupId, String consumerInstanceId) throws ExecutionException, InterruptedException { Properties props = new Properties(); props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); try (Admin admin = Admin.create(props)) { return admin.describeConsumerGroups(Set.of(groupId)).all().get() .get(groupId).members().stream() // 匹配实例ID,根据实际场景可匹配client.id或者consumer.id .filter(member -> member.consumerId().equals(consumerInstanceId)) .flatMap(member -> member.assignment().topicPartitions().stream()) .map(TopicPartition::topic) .collect(Collectors.toSet()); } } }
Restful 服务实现方案
基于Spring Boot封装上述逻辑为HTTP接口即可,示例如下:
1. Controller层代码
import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import java.util.Set; @RestController @RequestMapping("/kafka/consumer") public class KafkaConsumerController { @GetMapping("/topics/group/{groupId}") public Set<String> getTopicsByGroupId(@PathVariable String groupId) throws Exception { return ConsumerTopicQuery.getTopicsByConsumerGroup(groupId); } @GetMapping("/topics/instance/{groupId}/{instanceId}") public Set<String> getTopicsByInstanceId(@PathVariable String groupId, @PathVariable String instanceId) throws Exception { return ConsumerTopicQuery.getTopicsByConsumerInstanceId(groupId, instanceId); } }
2. 接口调用示例
- 查询消费组关联Topic:
GET http://你的服务地址/kafka/consumer/topics/group/consumer1 - 查询单个实例关联Topic:
GET http://你的服务地址/kafka/consumer/topics/instance/consumer_group1/consumer_instance_123
注意事项
- AdminClient的版本尽量和Kafka集群版本保持一致,避免兼容性问题
- 开启ACL的Kafka集群需要给AdminClient配置的账号授予
DESCRIBE权限,权限范围覆盖消费者组即可 - 高并发查询场景下可以对结果做30s-1min的本地缓存,减少对Broker的请求压力
内容的提问来源于stack exchange,提问作者VKR
相关产品推荐
相关产品推荐

