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

如何用segmentio/kafka-go编程获取Kafka主题的消费者组列表

使用segmentio/kafka-go获取订阅指定主题的消费者组列表

完全可以通过编程实现类似kafka-consumer-groups.sh的功能,segmentio/kafka-go支持Kafka的消费者组元数据查询协议。你当前的代码仅获取了主题分区和基础元数据,没有调用消费者组相关的查询接口,所以无法得到目标信息。

实现步骤与代码示例

以下是完整的实现代码,可获取订阅指定主题的消费者组及其详细信息:

package main

import (
	"context"
	"fmt"
	"github.com/segmentio/kafka-go"
)

func main() {
	kafkaBroker := "localhost:9092"
	targetTopic := "topic.test"

	// 创建Kafka连接
	conn, err := kafka.Dial("tcp", kafkaBroker)
	if err != nil {
		panic(fmt.Sprintf("连接Kafka失败: %v", err))
	}
	defer conn.Close()

	// 1. 获取所有消费者组列表
	groups, err := conn.ListGroups(context.Background())
	if err != nil {
		panic(fmt.Sprintf("获取消费者组列表失败: %v", err))
	}

	// 2. 遍历每个组,检查是否订阅了目标主题
	for _, group := range groups.Groups {
		// 描述该消费者组的详细信息
		groupDesc, err := conn.DescribeGroups(context.Background(), []string{group.GroupID})
		if err != nil {
			fmt.Printf("描述消费者组%s失败: %v\n", group.GroupID, err)
			continue
		}

		// 检查组是否订阅了目标主题
		for _, desc := range groupDesc.Groups {
			for _, member := range desc.Members {
				for _, assignment := range member.Assignment.Topics {
					if assignment == targetTopic {
						// 输出该消费者组的相关信息
						fmt.Printf("消费者组ID: %s\n", desc.GroupID)
						fmt.Printf("消费者ID: %s\n", member.MemberID)
						fmt.Printf("客户端ID: %s\n", member.ClientID)
						fmt.Printf("主机地址: %s\n", member.ClientHost)
						// 获取偏移量与延迟信息
						offsets, err := conn.FetchOffsets(context.Background(), kafka.FetchOffsetsRequest{
							GroupID: desc.GroupID,
							Topics: []kafka.FetchOffsetsTopic{
								{Topic: targetTopic},
							},
						})
						if err == nil {
							for _, tp := range offsets.Topics {
								for _, p := range tp.Partitions {
									fmt.Printf("主题%s 分区%d: 当前偏移量%d, 日志末尾偏移量%d, 延迟%d\n",
										tp.Topic, p.Partition, p.Offset, p.Offset+p.Lag, p.Lag)
								}
							}
						}
						fmt.Println("---")
					}
				}
			}
		}
	}
}

代码关键接口说明

  • ListGroups:获取Kafka集群中所有消费者组的ID列表
  • DescribeGroups:查询指定消费者组的详细信息,包括成员、订阅主题、分配的分区等
  • FetchOffsets:获取消费者组在指定主题分区上的当前偏移量、日志末尾偏移量及延迟数据,对应kafka-consumer-groups.sh输出的核心字段

关于Kafka Proxy API的疑问

你提到的Kafka Proxy API(通常指Confluent REST Proxy)是Confluent提供的RESTful交互接口,方便非Go/Java等语言通过HTTP操作Kafka,但它不是Confluent专属Kafka实现的必需组件。原生Apache Kafka使用二进制协议通信,segmentio/kafka-go直接基于该标准协议开发,无需依赖Confluent的任何组件,只要你的Kafka集群是标准Apache Kafka或兼容实现(比如Confluent Platform),都可以正常使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:27:41