如何使用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
相关产品推荐
相关产品推荐

