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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 03:56:03