如何使用Java获取特定Kafka Topic的生产者详情
使用Java获取Kafka特定Topic的生产者与消费者详情
获取特定Topic的生产者相关信息
Kafka本身没有提供直接枚举所有生产者的API,但可以通过两种方式获取生产者的关键信息:应用内生产者的本地指标和结合Broker元数据与监控指标追踪全局生产者。
1. 获取应用内自身生产者的详情
如果是你自己的Java应用中的Kafka生产者,可以直接通过KafkaProducer的metrics()方法获取生产速率、错误数等核心指标:
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Map; import java.util.Properties; public class LocalProducerMetrics { public static void main(String[] args) { Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker-1:9092,kafka-broker-2:9092"); producerProps.put(ProducerConfig.CLIENT_ID_CONFIG, "order-service-producer"); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); try (KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps)) { // 遍历并筛选生产相关指标 Map<String, ? extends org.apache.kafka.common.Metric> metrics = producer.metrics(); for (Map.Entry<String, ? extends org.apache.kafka.common.Metric> metricEntry : metrics.entrySet()) { String metricName = metricEntry.getKey(); // 筛选发送速率、错误率、延迟等关键指标 if (metricName.contains("record-send-rate") || metricName.contains("error-rate") || metricName.contains("record-latency-avg")) { System.out.printf("指标: %s, 数值: %s%n", metricName, metricEntry.getValue().metricValue()); } } } } }
2. 追踪全局范围内的生产者(结合Broker元数据)
要获取所有向目标Topic发送消息的生产者,需要依赖Broker的JMX指标或监控系统,同时可以通过AdminClient获取Topic的分区元数据,结合Broker端的producer-metrics指标(按客户端ID分组)来定位生产者:
首先用AdminClient获取目标Topic的基础元数据:
import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.DescribeTopicsResult; import org.apache.kafka.clients.admin.TopicDescription; import java.util.Collections; import java.util.Properties; import java.util.concurrent.ExecutionException; public class TopicMetadataFetcher { public static void main(String[] args) throws ExecutionException, InterruptedException { Properties adminProps = new Properties(); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker-1:9092,kafka-broker-2:9092"); try (AdminClient adminClient = AdminClient.create(adminProps)) { String targetTopic = "user-behaviors"; DescribeTopicsResult topicResult = adminClient.describeTopics(Collections.singletonList(targetTopic)); TopicDescription topicDesc = topicResult.values().get(targetTopic).get(); System.out.printf("Topic: %s, 分区数: %d%n", topicDesc.name(), topicDesc.partitions().size()); // 每个分区的Leader信息可关联到对应Broker的生产者指标 topicDesc.partitions().forEach(partition -> System.out.printf("分区ID: %d, Leader Broker: %d%n", partition.partition(), partition.leader().id()) ); } } }
之后可以通过Broker的JMX指标kafka.producer:type=producer-metrics,client-id=*来获取对应客户端ID的生产者的详细统计数据。
获取特定Topic的消费者详情
通过Kafka的AdminClient API,可以直接获取消费目标Topic的消费者群组、消费偏移量、客户端ID和主机信息:
import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.ConsumerGroupListing; import org.apache.kafka.clients.admin.ListConsumerGroupsResult; import org.apache.kafka.clients.admin.ListConsumerGroupOffsetsResult; import org.apache.kafka.clients.admin.OffsetDescription; import java.util.Properties; import java.util.concurrent.ExecutionException; public class TopicConsumerDetails { public static void main(String[] args) throws ExecutionException, InterruptedException { Properties adminProps = new Properties(); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker-1:9092,kafka-broker-2:9092"); String targetTopic = "user-behaviors"; try (AdminClient adminClient = AdminClient.create(adminProps)) { // 遍历所有消费者群组 ListConsumerGroupsResult groupsResult = adminClient.listConsumerGroups(); for (ConsumerGroupListing group : groupsResult.all().get()) { // 获取该群组的所有消费偏移量 ListConsumerGroupOffsetsResult offsetsResult = adminClient.listConsumerGroupOffsets(group.groupId()); offsetsResult.partitionsToOffsetAndMetadata().get().forEach((tp, offsetDesc) -> { // 筛选目标Topic的消费信息 if (tp.topic().equals(targetTopic)) { System.out.printf("=== 消费者群组: %s ===%n", group.groupId()); System.out.printf("Topic分区: %d%n", tp.partition()); System.out.printf("当前消费偏移量: %d%n", offsetDesc.offset()); System.out.printf("消费者客户端ID: %s%n", offsetDesc.consumerId()); System.out.printf("消费者主机: %s%n", offsetDesc.host()); } }); } } } }
注意事项
- 版本兼容:确保
kafka-clients依赖版本与Kafka集群版本一致,避免API不兼容问题,例如使用org.apache.kafka:kafka-clients:2.8.2对应Kafka 2.8.x集群。 - 权限要求:
AdminClient需要具备DescribeGroups、DescribeConsumerGroups和DescribeTopics等权限,否则无法获取完整信息。 - 全局生产者限制:Kafka没有集中存储所有生产者的列表,全局生产者信息只能通过Broker端的监控指标间接获取,无法通过API直接枚举。
内容的提问来源于stack exchange,提问作者Pramod Kumar
相关产品推荐
相关产品推荐

