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
相关产品推荐
相关产品推荐

