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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 06:57:10