如何用Golang Sarama客户端检测Kafka主题新增分区?
使用Sarama检测Kafka主题新增分区的几种方法
当然可以通过Sarama检测Kafka主题的新增分区,除了你想到的定期轮询主题信息,还有更高效的事件驱动方式,以下是具体实现思路:
1. 定期拉取主题元数据(你已想到的方案)
通过Sarama的Client实例调用DescribeTopics或Topics方法,定时获取指定主题的元数据,对比历史记录的分区列表/数量,就能识别新增分区。这种方式实现简单,适合对实时性要求不高的场景。
示例代码片段:
package main import ( "fmt" "time" "github.com/Shopify/sarama" ) func main() { config := sarama.NewConfig() client, err := sarama.NewClient([]string{"kafka-broker:9092"}, config) if err != nil { panic(err) } defer client.Close() topic := "your-topic" prevPartitions := make(map[int32]struct{}) // 初始化获取初始分区 topics, err := client.DescribeTopics([]string{topic}) if err != nil { panic(err) } for _, p := range topics[0].Partitions { prevPartitions[p.ID] = struct{}{} } // 定时轮询 ticker := time.NewTicker(30 * time.Second) defer ticker.Stop() for range ticker.C { topics, err := client.DescribeTopics([]string{topic}) if err != nil { fmt.Printf("Failed to describe topic: %v\n", err) continue } currentPartitions := make(map[int32]struct{}) for _, p := range topics[0].Partitions { currentPartitions[p.ID] = struct{}{} } // 检测新增分区 for pID := range currentPartitions { if _, exists := prevPartitions[pID]; !exists { fmt.Printf("New partition detected: %d\n", pID) } } prevPartitions = currentPartitions } }
2. 利用ConsumerGroup的Rebalance事件(更高效的事件驱动方案)
当Kafka主题新增分区时,对应的ConsumerGroup会触发Rebalance操作。Sarama的ConsumerGroupHandler接口提供了Setup和Cleanup方法,在Rebalance完成后,Setup会被调用,此时可以获取到当前分配给消费者的所有分区,通过对比历史分配记录,就能实时感知新增分区。
这种方式不需要主动轮询,是事件驱动的,实时性更高,适合需要及时处理新增分区的业务场景。
示例代码片段:
package main import ( "context" "fmt" "sync" "github.com/Shopify/sarama" ) type MyConsumerGroupHandler struct { assignedPartitions map[string][]int32 mu sync.Mutex } func (h *MyConsumerGroupHandler) Setup(sess sarama.ConsumerGroupSession) error { h.mu.Lock() defer h.mu.Unlock() current := make(map[string][]int32) for topic, partitions := range sess.Claims() { current[topic] = partitions } // 对比检测新增分区 for topic, parts := range current { prevParts, exists := h.assignedPartitions[topic] if !exists { fmt.Printf("Topic %s: initial partitions assigned: %v\n", topic, parts) h.assignedPartitions[topic] = parts continue } prevMap := make(map[int32]struct{}) for _, p := range prevParts { prevMap[p] = struct{}{} } for _, p := range parts { if _, exists := prevMap[p]; !exists { fmt.Printf("Topic %s: new partition detected: %d\n", topic, p) } } h.assignedPartitions[topic] = parts } return nil } func (h *MyConsumerGroupHandler) Cleanup(sess sarama.ConsumerGroupSession) error { return nil } func (h *MyConsumerGroupHandler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for msg := range claim.Messages() { // 处理消息逻辑 sess.MarkMessage(msg, "") } return nil } func main() { config := sarama.NewConfig() config.Consumer.Return.Errors = true config.Version = sarama.V2_0_0_0 // 根据你的Kafka版本调整 group, err := sarama.NewConsumerGroup([]string{"kafka-broker:9092"}, "your-group-id", config) if err != nil { panic(err) } defer group.Close() handler := &MyConsumerGroupHandler{ assignedPartitions: make(map[string][]int32), } wg := &sync.WaitGroup{} wg.Add(1) ctx := context.Background() go func() { defer wg.Done() for { err := group.Consume(ctx, []string{"your-topic"}, handler) if err != nil { fmt.Printf("Consumer group error: %v\n", err) } } }() wg.Wait() }
两种方案对比
- 定期轮询:实现简单,无需依赖ConsumerGroup,但存在延迟,适合实时性要求低的场景。
- Rebalance事件驱动:实时性高,无轮询开销,但需要基于ConsumerGroup架构,适合需要及时处理新增分区的业务。
内容的提问来源于stack exchange,提问作者A Beginner
相关产品推荐
相关产品推荐

