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

使用kafka-go与RoundRobin负载均衡器时数据始终写入分区0的问题

问题分析与解决方案

你的代码里所有消息都写入分区0,核心原因有两个:

1. 每条消息属于不同Topic且仅写入一次

kafka-go的RoundRobin负载均衡器是按Topic独立维护轮询状态的,每个Topic首次写入时都会从分区0开始分配。你现在给两个不同的Topic各写一条消息,自然每条都会落到对应Topic的分区0上。

2. 未验证同一Topic下的多消息写入

如果要测试RoundRobin的轮询效果,需要往同一个Topic写入多条消息——不管是批量写入还是多次单条写入,轮询逻辑都会正常工作。

另外补充一个规范建议:初始化Writer时最好指定Topic字段,避免每条消息都重复写Topic名,代码更简洁。

修正后的示例代码

import "fmt"

func main() {
    var logger = logger.Logger()

    w := kafka.Writer{
        Addr:     kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"),
        Balancer: &kafka.RoundRobin{},
        Topic:    "first-topic", // 指定默认Topic,消息中可省略Topic字段
    }

    // 同一Topic下的多条消息,批量写入
    messages := []kafka.Message{
        {Key: []byte("test"), Value: []byte("value")},
        {Key: []byte("test2"), Value: []byte("value2")},
        {Key: []byte("test3"), Value: []byte("value3")},
    }
    if err := w.WriteMessages(context.Background(), messages...); err != nil {
        logger.Fatal().Msgf("写入消息失败: %v", err)
    }

    // 循环写入单条消息,验证轮询效果
    for i := 0; i < 5; i++ {
        msg := kafka.Message{
            Key:   []byte(fmt.Sprintf("loop-test-%d", i)),
            Value: []byte(fmt.Sprintf("loop-value-%d", i)),
        }
        if err := w.WriteMessages(context.Background(), msg); err != nil {
            logger.Error().Msgf("写入单条消息失败: %v", err)
        }
    }

    if err := w.Close(); err != nil {
        logger.Fatal().Msgf("关闭Writer失败: %v", err)
    }
}

额外注意事项

  • 确保你的目标Topic创建时指定了多个分区(比如使用命令kafka-topics.sh --create --topic first-topic --bootstrap-server localhost:9092 --partitions 3),否则只有分区0可用,轮询逻辑无法生效。
  • RoundRobin的轮询计数器是按Topic隔离的,不同Topic的分区选择互不干扰。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 13:05:12