我正在尝试使用 Golang 客户端测试生产者将消息写入 kafka 集群上的主题。package mainimport ( "fmt" "gopkg.in/confluentinc/confluent-kafka-go.v1/kafka")func main() { p, err := kafka.NewProducer(&kafka.ConfigMap{"bootstrap.servers":"localhost"}) if err != nil { panic(err) } defer p.Close() // Delivery report handler for produced messages go func() { for e := range p.Events() { switch ev := e.(type) { case *kafka.Message: if ev.TopicPartition.Error != nil { fmt.Printf("Delivery failed: %v\n", ev.TopicPartition) } else { fmt.Printf("Delivered message to %v\n", ev.TopicPartition) } } } }() // Produce messages to topic (asynchronously) topic := "test" for _, word := range []string{"test message"} { p.Produce(&kafka.Message{ TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: kafka.PartitionAny}, Value: []byte(word), }, nil) } // Wait for message deliveries before shutting down p.Flush(15 * 1000)}我在控制台上收到消息-消费者没有问题。然后,我尝试做同样的事情,只是使用我的远程 kafka 集群主题(注意我也尝试过不使用字符串中的端口):p, err := kafka.NewProducer(&kafka.ConfigMap{"bootstrap.servers":"HOSTNAME.amazonaws.com:9092,HOSTNAME2.amazonaws.com:9092,HOSTNAME3.amazonaws.com:9092"})它打印以下错误:Delivery failed: test[0]@end(Broker: Not enough in-sync replicas)不过控制台生产者没有任何问题:./bin/kafka-console-producer.sh --broker-list HOSTNAME.amazonaws.com:9092,HOSTNAME2.amazonaws.com:9092,HOSTNAME3.amazonaws.com:9092 --topic test>proving that this works控制台消费者收到它:bin/kafka-console-consumer.sh --bootstrap-server HOSTNAME.amazonaws.com:9092,HOSTNAME2.amazonaws.com:9092,HOSTNAME3.amazonaws.com:9092 --topic test --from-beginning proving that this works我做的最后一件事是检查该主题有多少个同步副本。如果我没看错的话,最小值应该是 2,而且有 3 个。我还有什么可以研究的想法吗?
2 回答
莫回无
TA贡献1865条经验 获得超7个赞
您有min.insync.replicas=2
,但该主题只有一个副本。
我相信控制台制作人只将该属性设置为 1
有 3 个
其实只有一个。这是代理 ID 3。如果实际上有三个副本,您会看到总共三个单独的数字作为 ISR
- 2 回答
- 0 关注
- 143 浏览
添加回答
举报
0/150
提交
取消