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

MacOS下Go Kafka程序运行报错:Not Leader For Partition

解决Kafka客户端报错"Not Leader For Partition"的方案

问题描述

运行Go语言Kafka读写程序时出现运行时错误:

2023/06/17 20:46:11 Failed to read message: [6] Not Leader For Partition: the client attempted to send messages to a replica that is not the leader for some partition, the client's metadata are likely out of date

已通过brew确认Kafka、ZooKeeper服务均正常启动。

解决步骤

1. 手动创建目标Topic

自动创建Topic可能引发元数据同步延迟,先手动创建指定Topic:

kafka-topics --create --topic my-topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1

2. 检查Kafka Broker配置

brew安装的Kafka配置文件路径为/usr/local/etc/kafka/server.properties,需确保以下配置正确:

  • 确认listeners设置为PLAINTEXT://localhost:9092,保证Broker监听本地地址
  • 确保advertised.listeners与listeners配置一致,避免客户端获取错误的Broker地址

3. 修改Go客户端配置

调整客户端参数,避免元数据不一致问题:

  • 写入器禁用自动创建Topic,增加重试次数
  • 读取器指定分区,设置合理的等待时间

修改后的完整代码:

package main

import (
        "context"
        "fmt"
        "log"
        "os"
        "os/signal"
        "sync"

        "github.com/segmentio/kafka-go"
)

const (
        brokerAddress = "localhost:9092"
        topic         = "my-topic"
)

func main() {
        // 创建Kafka写入器,禁用自动创建topic并增加重试
        writer := kafka.NewWriter(kafka.WriterConfig{
                Brokers:               []string{brokerAddress},
                Topic:                 topic,
                Balancer:              &kafka.LeastBytes{},
                AllowAutoTopicCreation: false, // 禁用自动创建topic,避免元数据不一致
                MaxAttempts:           3,      // 增加重试次数
        })

        // 创建Kafka读取器,指定分区并设置等待时间
        reader := kafka.NewReader(kafka.ReaderConfig{
                Brokers:   []string{brokerAddress},
                Topic:     topic,
                Partition: 0, // 指定分区,避免元数据未同步时的分区查找问题
                MaxWait:   100, // 设置最大等待时间,加快元数据刷新
        })

        // 启动消费协程
        go func() {
                for {
                        message, err := reader.ReadMessage(context.Background())
                        if err != nil {
                                log.Println("读取消息失败:", err)
                                continue
                        }
                        fmt.Printf("收到消息: %s\n", message.Value)
                }
        }()

        // 生产测试消息
        err := writer.WriteMessages(context.Background(),
                kafka.Message{
                        Key:   []byte("key"),
                        Value: []byte("Hello, Kafka!"),
                },
        )
        if err != nil {
                log.Println("生产消息失败:", err)
        }

        // 等待中断信号退出
        wg := sync.WaitGroup{}
        wg.Add(1)
        go func() {
                defer wg.Done()
                c := make(chan os.Signal, 1)
                signal.Notify(c, os.Interrupt)
                <-c
        }()

        wg.Wait()

        // 关闭资源
        if err := writer.Close(); err != nil {
                log.Println("关闭写入器失败:", err)
        }
        if err := reader.Close(); err != nil {
                log.Println("关闭读取器失败:", err)
        }
}

4. 重启Kafka服务

应用配置修改后,重启Kafka:

brew services restart kafka

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 13:33:14