使用kafka-go官方示例代码消费消息时关闭batch触发请求超时错误的原因咨询
kafka-go官方示例代码消费消息时关闭batch触发请求超时错误的原因咨询
我刚在一台全新的机器上本地搭建了Kafka环境,这是我第一次使用kafka-go库。我直接复制了官方文档里的生产消息和消费消息两段示例代码,拼接成了完整的程序,代码如下:
package main import ( "context" "fmt" "log" "time" "github.com/segmentio/kafka-go" ) func main() { ProduceMessages() ConsumeMessages() } func ProduceMessages() { // to produce messages topic := "my-topic" partition := 0 conn, err := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", topic, partition) if err != nil { log.Fatal("failed to dial leader:", err) } conn.SetWriteDeadline(time.Now().Add(10 * time.Second)) _, err = conn.WriteMessages( kafka.Message{Value: []byte("one!")}, kafka.Message{Value: []byte("two!")}, kafka.Message{Value: []byte("three!")}, ) if err != nil { log.Fatal("failed to write messages:", err) } if err := conn.Close(); err != nil { log.Fatal("failed to close writer:", err) } } func ConsumeMessages() { // to consume messages topic := "my-topic" partition := 0 conn, err := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", topic, partition) if err != nil { log.Fatal("failed to dial leader:", err) } conn.SetReadDeadline(time.Now().Add(10 * time.Second)) batch := conn.ReadBatch(10e3, 1e6) // fetch 10KB min, 1MB max b := make([]byte, 10e3) // 10KB max per message for { n, err := batch.Read(b) if err != nil { break } fmt.Println(string(b[:n])) } if err := batch.Close(); err != nil { log.Fatal("failed to close batch:", err) } if err := conn.Close(); err != nil { log.Fatal("failed to close connection:", err) } }
运行这段代码后,消息确实能被正常消费并打印出来:
one!
two!
three!
但等待几秒后,程序总会抛出一个超时错误:
2022/12/10 10:35:50 failed to close batch:[7] Request Timed Out: the request exceeded the user-specified time limit in the request
exit status 1
我想请教各位,为什么关闭batch的时候会触发这个请求超时的问题?
PS:我无法添加kafka-go这个新标签,只能使用已有的confluent-kafka-go标签。
备注:内容来源于stack exchange,提问作者Mohamad-Jaafar NEHME
相关产品推荐
相关产品推荐

