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

如何使用Confluent-Go Admin Client获取Kafka Topic配置及副本信息?

问题解答

1. 获取Topic完整配置(含retention.ms、min.isr等)

Confluent Go Admin API的Metadata接口仅返回基础元数据,要获取完整的Topic配置,需使用DescribeTopics方法。该方法会返回指定Topic的所有配置项(包括自定义配置与默认配置),对应TopicDescription结构体中的Configs字段。

2. 获取Observers与Offline副本信息

  • Observers:属于Confluent Kafka的扩展特性,可通过PartitionDescription结构体中的Observers字段直接获取对应副本列表。
  • Offline副本:PartitionDescription结构体提供了OfflineReplicas字段,可直接拿到处于离线状态的副本集合。

示例代码(模拟kafka-topics输出格式)

package main

import (
	"context"
	"fmt"
	"strings"

	"github.com/confluentinc/confluent-kafka-go/v2/kafka"
)

func main() {
	// 初始化Admin客户端
	adminClient, err := kafka.NewAdminClient(&kafka.ConfigMap{
		"bootstrap.servers": "localhost:9092",
	})
	if err != nil {
		fmt.Printf("Failed to create admin client: %v\n", err)
		return
	}
	defer adminClient.Close()

	// 指定要查询的Topic
	topicNames := []string{"test-topic"}
	describeOpts := kafka.DescribeTopicsOptions{}

	// 调用DescribeTopics获取完整信息
	topicDescriptions, err := adminClient.DescribeTopics(context.Background(), topicNames, describeOpts)
	if err != nil {
		fmt.Printf("Failed to describe topics: %v\n", err)
		return
	}

	// 格式化输出,模拟kafka-topics的样式
	for _, td := range topicDescriptions {
		// 拼接配置字符串
		var configStrs []string
		for _, cfg := range td.Configs {
			configStrs = append(configStrs, fmt.Sprintf("%s=%s", cfg.Name, cfg.Value))
		}
		configLine := strings.Join(configStrs, ",")

		// 输出Topic级信息
		fmt.Printf("Topic: %-20s TopicId: %s PartitionCount: %d       ReplicationFactor: %d    Configs: %s\n",
			td.Name, td.TopicID, len(td.Partitions), td.ReplicationFactor, configLine)

		// 输出每个分区的详细信息
		for _, pd := range td.Partitions {
			// 转换副本列表为字符串
			replicasStr := strings.Join(strings.Fields(fmt.Sprint(pd.Replicas)), ",")
			isrStr := strings.Join(strings.Fields(fmt.Sprint(pd.Isr)), ",")
			offlineStr := strings.Join(strings.Fields(fmt.Sprint(pd.OfflineReplicas)), ",")
			observersStr := strings.Join(strings.Fields(fmt.Sprint(pd.Observers)), ",")

			// 处理空值显示
			if offlineStr == "" {
				offlineStr = "        "
			}
			if observersStr == "" {
				observersStr = "        "
			}

			fmt.Printf("\tTopic: %-20s Partition: %d    Leader: %d       Replicas: %s Isr: %s        Offline: %s Observers: %s\n",
				td.Name, pd.Partition, pd.Leader, replicasStr, isrStr, offlineStr, observersStr)
		}
	}
}

输出说明

运行上述代码后,输出格式会与kafka-topics --describe的结果一致,包含Topic ID、分区数、副本因子、所有配置项,以及每个分区的Leader、副本列表、ISR、离线副本和Observer副本信息。

内容的提问来源于stack exchange,提问作者Neeraj

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 21:40:04