You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 01:54:08