segmentio/kafka-go

GitHub: segmentio/kafka-go

一个纯 Go 实现的 Kafka 客户端库,通过贴合 Go 标准库理念的高底层 API 解决了现有 Go Kafka 客户端难用、依赖重或不支持 context 的问题。

Stars: 8592 | Forks: 853

# kafka-go [![CircleCI](https://circleci.com/gh/segmentio/kafka-go.svg?style=shield)](https://circleci.com/gh/segmentio/kafka-go) [![Go Report Card](https://goreportcard.com/badge/github.com/segmentio/kafka-go)](https://goreportcard.com/report/github.com/segmentio/kafka-go) [![GoDoc](https://godoc.org/github.com/segmentio/kafka-go?status.svg)](https://godoc.org/github.com/segmentio/kafka-go) ## 动机 在 Segment,我们非常依赖 Go 和 Kafka。然而,截至撰写本文时,Go 中的 Kafka 客户端库现状并不理想。可用的选项有: - [sarama](https://github.com/Shopify/sarama),这是目前最流行的库,但使用起来相当困难。它的文档匮乏,API 暴露了 Kafka 协议的底层概念,并且不支持诸如 [context](https://golang.org/pkg/context/) 等 Go 的新特性。它还将所有值作为指针传递,这会导致大量的动态内存分配、更频繁的垃圾回收以及更高的内存占用。 - [confluent-kafka-go](https://github.com/confluentinc/confluent-kafka-go) 是一个基于 cgo 的对 [librdkafka](https://github.com/edenhill/librdkafka) 的封装,这意味着它为所有使用该包的 Go 代码引入了对 C 依赖库的依赖。它的文档比 sarama 好得多,但仍然缺乏对 Go context 的支持。 - [goka](https://github.com/lovoo/goka) 是一个较新的 Go Kafka 客户端,专注于特定的使用模式。它提供了将 Kafka 用作服务间消息传递总线的抽象,而不是有序的事件日志,但这在 Segment 并不是我们典型的 Kafka 使用场景。该包在所有与 Kafka 的交互中也依赖于 sarama。 这就是 `kafka-go` 发挥作用的地方。它提供了与 Kafka 交互的底层和高层 API,镜像了 Go 标准库的概念并实现了其接口,从而使其易于使用并与现有软件集成。 #### 注意: 为了更好地遵守我们新采用的《行为准则》,kafka-go 项目已将默认分支重命名为 `main`。有关我们行为准则的完整详细信息,请参阅[此](./CODE_OF_CONDUCT.md)文档。 ## Kafka 版本 `kafka-go` 目前针对 Kafka 0.10.1.0 到 2.7.1 版本进行测试。 虽然它也应该与更高版本兼容,但 Kafka API 中可用的新特性可能尚未在客户端中实现。 ## Go 版本 `kafka-go` 需要 Go 1.15 或更高版本。 ## 连接 [![GoDoc](https://godoc.org/github.com/segmentio/kafka-go?status.svg)](https://godoc.org/github.com/segmentio/kafka-go#Conn) `Conn` 类型是 `kafka-go` 包的核心。它封装了原始网络连接,以向 Kafka 服务器提供底层 API。 以下是一些展示连接对象典型用法的示例: ``` // 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) } ``` ``` // 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) } ``` ### 创建 Topic 默认情况下,Kafka 开启了 `auto.create.topics.enable='true'`(在 bitnami/kafka 的 Kafka docker 镜像中为 `KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE='true'`)。如果此值设置为 `'true'`,那么作为 `kafka.DialLeader` 的副作用,topic 将被创建,如下所示: ``` // to create topics when auto.create.topics.enable='true' conn, err := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", "my-topic", 0) if err != nil { panic(err.Error()) } ``` 如果 `auto.create.topics.enable='false'`,则你需要像这样显式地创建 topic: ``` // to create topics when auto.create.topics.enable='false' topic := "my-topic" conn, err := kafka.Dial("tcp", "localhost:9092") if err != nil { panic(err.Error()) } defer conn.Close() controller, err := conn.Controller() if err != nil { panic(err.Error()) } var controllerConn *kafka.Conn controllerConn, err = kafka.Dial("tcp", net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port))) if err != nil { panic(err.Error()) } defer controllerConn.Close() topicConfigs := []kafka.TopicConfig{ { Topic: topic, NumPartitions: 1, ReplicationFactor: 1, }, } err = controllerConn.CreateTopics(topicConfigs...) if err != nil { panic(err.Error()) } ``` ### 通过非 leader 连接来连接到 leader ``` // to connect to the kafka leader via an existing non-leader connection rather than using DialLeader conn, err := kafka.Dial("tcp", "localhost:9092") if err != nil { panic(err.Error()) } defer conn.Close() controller, err := conn.Controller() if err != nil { panic(err.Error()) } var connLeader *kafka.Conn connLeader, err = kafka.Dial("tcp", net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port))) if err != nil { panic(err.Error()) } defer connLeader.Close() ``` ### 列出 topic ``` conn, err := kafka.Dial("tcp", "localhost:9092") if err != nil { panic(err.Error()) } defer conn.Close() partitions, err := conn.ReadPartitions() if err != nil { panic(err.Error()) } m := map[string]struct{}{} for _, p := range partitions { m[p.Topic] = struct{}{} } for k := range m { fmt.Println(k) } ``` 因为是底层的,所以 `Conn` 类型成为了构建更高层次抽象(例如 `Reader`)的绝佳基础。 ## Reader [![GoDoc](https://godoc.org/github.com/segmentio/kafka-go?status.svg)](https://godoc.org/github.com/segmentio/kafka-go#Reader) `Reader` 是 `kafka-go` 包暴露的另一个概念,旨在简化从单个 topic-partition 对进行消费的典型使用场景的实现。 `Reader` 还会自动处理重连和 offset 管理,并暴露了支持使用 Go context 进行异步取消和超时的 API。 请注意,在进程退出时对 `Reader` 调用 `Close()` 非常重要。Kafka 服务器需要优雅断开连接,以阻止其继续尝试向已连接的客户端发送消息。如果进程被 SIGINT(在 shell 中按 ctrl-c)或 SIGTERM(如 docker stop 或 kubernetes 重启)终止,给定的示例将不会调用 `Close()`。这可能会导致同一个 topic 上的新 reader 连接时出现延迟(例如,启动新进程或运行新容器)。请使用 `signal.Notify` 处理程序在进程关闭时关闭 reader。 ``` // make a new reader that consumes from topic-A, partition 0, at offset 42 r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092","localhost:9093", "localhost:9094"}, Topic: "topic-A", Partition: 0, MaxBytes: 10e6, // 10MB }) r.SetOffset(42) for { m, err := r.ReadMessage(context.Background()) if err != nil { break } fmt.Printf("message at offset %d: %s = %s\n", m.Offset, string(m.Key), string(m.Value)) } if err := r.Close(); err != nil { log.Fatal("failed to close reader:", err) } ``` ### 消费者组 ```kafka-go``` 还支持 Kafka 消费者组,包括由 broker 管理的 offset。 要启用消费者组,只需在 ReaderConfig 中指定 GroupID 即可。 使用消费者组时,ReadMessage 会自动提交 offset。 ``` // make a new reader that consumes from topic-A r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"}, GroupID: "consumer-group-id", Topic: "topic-A", MaxBytes: 10e6, // 10MB }) for { m, err := r.ReadMessage(context.Background()) if err != nil { break } fmt.Printf("message at topic/partition/offset %v/%v/%v: %s = %s\n", m.Topic, m.Partition, m.Offset, string(m.Key), string(m.Value)) } if err := r.Close(); err != nil { log.Fatal("failed to close reader:", err) } ``` 使用消费者组时有诸多限制: * 当设置了 GroupID 时,```(*Reader).SetOffset``` 将返回错误 * 当设置了 GroupID 时,```(*Reader).Offset``` 将始终返回 ```-1``` * 当设置了 GroupID 时,```(*Reader).Lag``` 将始终返回 ```-1``` * 当设置了 GroupID 时,```(*Reader).ReadLag``` 将返回错误 * 当设置了 GroupID 时,```(*Reader).Stats``` 将返回 ```-1``` 的分区 ### 显式提交 ```kafka-go``` 还支持显式提交。你可以调用 ```FetchMessage``` 然后再调用 ```CommitMessages```,而不是调用 ```ReadMessage```。 ``` ctx := context.Background() for { m, err := r.FetchMessage(ctx) if err != nil { break } fmt.Printf("message at topic/partition/offset %v/%v/%v: %s = %s\n", m.Topic, m.Partition, m.Offset, string(m.Key), string(m.Value)) if err := r.CommitMessages(ctx, m); err != nil { log.Fatal("failed to commit messages:", err) } } ``` 在消费者组中提交消息时,对于给定的 topic/partition,具有最高 offset 的消息决定了该分区已提交的 offset 的值。例如,如果通过调用 `FetchMessage` 获取了单个分区中 offset 为 1、2 和 3 的消息,那么使用 offset 为 3 的消息调用 `CommitMessages` 也会导致提交该分区中 offset 为 1 和 2 的消息。 ### 管理提交 默认情况下,CommitMessages 会同步将 offset 提交到 Kafka。为了提高性能,你可以改为通过在 ReaderConfig 上设置 CommitInterval 来周期性地向 Kafka 提交 offset。 ``` // make a new reader that consumes from topic-A r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"}, GroupID: "consumer-group-id", Topic: "topic-A", MaxBytes: 10e6, // 10MB CommitInterval: time.Second, // flushes commits to Kafka every second }) ``` ## Writer [![GoDoc](https://godoc.org/github.com/segmentio/kafka-go?status.svg)](https://godoc.org/github.com/segmentio/kafka-go#Writer) 为了向 Kafka 生产消息,程序可以使用底层的 `Conn` API,但该包还提供了更高层次的 `Writer` 类型。这在大多数情况下更为适用,因为它提供了额外的功能: - 发生错误时自动重试和重连。 - 可配置的消息在可用分区间的分配方式。 - 向 Kafka 同步或异步写入消息。 - 使用 context 进行异步取消。 - 关闭时刷新挂起的消息以支持优雅关闭。 - 在发布消息前创建缺失的 topic。*注意!* 在版本 `v0.4.30` 之前,这是默认行为。 ``` // make a writer that produces to topic-A, using the least-bytes distribution w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Balancer: &kafka.LeastBytes{}, } err := w.WriteMessages(context.Background(), kafka.Message{ Key: []byte("Key-A"), Value: []byte("Hello World!"), }, kafka.Message{ Key: []byte("Key-B"), Value: []byte("One!"), }, kafka.Message{ Key: []byte("Key-C"), Value: []byte("Two!"), }, ) if err != nil { log.Fatal("failed to write messages:", err) } if err := w.Close(); err != nil { log.Fatal("failed to close writer:", err) } ``` ### 发布前创建缺失的 topic ``` // Make a writer that publishes messages to topic-A. // The topic will be created if it is missing. w := &Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", AllowAutoTopicCreation: true, } messages := []kafka.Message{ { Key: []byte("Key-A"), Value: []byte("Hello World!"), }, { Key: []byte("Key-B"), Value: []byte("One!"), }, { Key: []byte("Key-C"), Value: []byte("Two!"), }, } var err error const retries = 3 for i := 0; i < retries; i++ { ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() // attempt to create topic prior to publishing the message err = w.WriteMessages(ctx, messages...) if errors.Is(err, kafka.LeaderNotAvailable) || errors.Is(err, context.DeadlineExceeded) { time.Sleep(time.Millisecond * 250) continue } if err != nil { log.Fatalf("unexpected error %v", err) } break } if err := w.Close(); err != nil { log.Fatal("failed to close writer:", err) } ``` ### 写入多个 topic 通常,`WriterConfig.Topic` 用于初始化单 topic 的 writer。排除此特定配置后,你可以通过设置 `Message.Topic` 来按消息定义 topic。 ``` w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), // NOTE: When Topic is not defined here, each Message must define it instead. Balancer: &kafka.LeastBytes{}, } err := w.WriteMessages(context.Background(), // NOTE: Each Message has Topic defined, otherwise an error is returned. kafka.Message{ Topic: "topic-A", Key: []byte("Key-A"), Value: []byte("Hello World!"), }, kafka.Message{ Topic: "topic-B", Key: []byte("Key-B"), Value: []byte("One!"), }, kafka.Message{ Topic: "topic-C", Key: []byte("Key-C"), Value: []byte("Two!"), }, ) if err != nil { log.Fatal("failed to write messages:", err) } if err := w.Close(); err != nil { log.Fatal("failed to close writer:", err) } ``` **注意:** 这 2 种模式是互斥的,如果你设置了 `Writer.Topic`,就不应在你写入的消息上显式定义 `Message.Topic`。反之,当你没有为 writer 定义 topic 时也一样。如果 `Writer` 检测到这种歧义,将会返回错误。 ### 与其他客户端的兼容性 #### Sarama 如果你正在从 Sarama 切换过来,并且需要/希望使用相同的消息分区算法,你可以使用 `kafka.Hash` balancer 或 `kafka.ReferenceHash` balancer: * `kafka.Hash` = `sarama.NewHashPartitioner` * `kafka.ReferenceHash` = `sarama.NewReferenceHashPartitioner` `kafka.Hash` 和 `kafka.ReferenceHash` balancer 会将消息路由到与上述两个 Sarama 分区器将其路由到的相同分区。 ``` w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Balancer: &kafka.Hash{}, } ``` #### librdkafka 和 confluent-kafka-go 使用 ```kafka.CRC32Balancer``` balancer 以获得与 librdkafka 默认的 ```consistent_random``` 分区策略相同的行为。 ``` w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Balancer: kafka.CRC32Balancer{}, } ``` #### Java 使用 ```kafka.Murmur2Balancer``` balancer 以获得与标准 Java 客户端默认分区器相同的行为。注意:Java 类允许你直接指定分区,而这在这里是不允许的。 ``` w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Balancer: kafka.Murmur2Balancer{}, } ``` ### 压缩 可以通过设置 `Compression` 字段在 `Writer` 上启用压缩: ``` w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Compression: kafka.Snappy, } ``` `Reader` 将通过检查消息属性来确定消费的消息是否已压缩。但是,必须导入所有预期 codec 的包,以便它们能被正确加载。 _注意:在 0.4 之前的版本中,程序必须导入压缩包才能安装 codec 并支持从 kafka 读取压缩消息。现在情况已不再如此,导入压缩包现在已是无操作(no-ops)。_ ### 连接 ``` dialer := &kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, TLS: &tls.Config{...tls config...}, } conn, err := dialer.DialContext(ctx, "tcp", "localhost:9093") ``` ### Reader ``` dialer := &kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, TLS: &tls.Config{...tls config...}, } r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"}, GroupID: "consumer-group-id", Topic: "topic-A", Dialer: dialer, }) ``` ### Writer 直接创建 Writer ``` w := kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Balancer: &kafka.Hash{}, Transport: &kafka.Transport{ TLS: &tls.Config{}, }, } ``` 使用 `kafka.NewWriter` ``` dialer := &kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, TLS: &tls.Config{...tls config...}, } w := kafka.NewWriter(kafka.WriterConfig{ Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"}, Topic: "topic-A", Balancer: &kafka.Hash{}, Dialer: dialer, }) ``` 请注意,`kafka.NewWriter` 和 `kafka.WriterConfig` 已被弃用,并将在未来的版本中删除。 ### SASL 认证类型 #### [Plain](https://godoc.org/github.com/segmentio/kafka-go/sasl/plain#Mechanism) ``` mechanism := plain.Mechanism{ Username: "username", Password: "password", } ``` #### [SCRAM](https://godoc.org/github.com/segmentio/kafka-go/sasl/scram#Mechanism) ``` mechanism, err := scram.Mechanism(scram.SHA512, "username", "password") if err != nil { panic(err) } ``` ### 连接 ``` mechanism, err := scram.Mechanism(scram.SHA512, "username", "password") if err != nil { panic(err) } dialer := &kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, SASLMechanism: mechanism, } conn, err := dialer.DialContext(ctx, "tcp", "localhost:9093") ``` ### Reader ``` mechanism, err := scram.Mechanism(scram.SHA512, "username", "password") if err != nil { panic(err) } dialer := &kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, SASLMechanism: mechanism, } r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092","localhost:9093", "localhost:9094"}, GroupID: "consumer-group-id", Topic: "topic-A", Dialer: dialer, }) ``` ### Writer ``` mechanism, err := scram.Mechanism(scram.SHA512, "username", "password") if err != nil { panic(err) } // Transports are responsible for managing connection pools and other resources, // it's generally best to create a few of these and share them across your // application. sharedTransport := &kafka.Transport{ SASL: mechanism, } w := kafka.Writer{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Topic: "topic-A", Balancer: &kafka.Hash{}, Transport: sharedTransport, } ``` ### 客户端 ``` mechanism, err := scram.Mechanism(scram.SHA512, "username", "password") if err != nil { panic(err) } // Transports are responsible for managing connection pools and other resources, // it's generally best to create a few of these and share them across your // application. sharedTransport := &kafka.Transport{ SASL: mechanism, } client := &kafka.Client{ Addr: kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"), Timeout: 10 * time.Second, Transport: sharedTransport, } ``` #### 读取时间范围内的所有消息 ``` startTime := time.Now().Add(-time.Hour) endTime := time.Now() batchSize := int(10e6) // 10MB r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"}, Topic: "my-topic1", Partition: 0, MaxBytes: batchSize, }) r.SetOffsetAt(context.Background(), startTime) for { m, err := r.ReadMessage(context.Background()) if err != nil { break } if m.Time.After(endTime) { break } // TODO: process message fmt.Printf("message at offset %d: %s = %s\n", m.Offset, string(m.Key), string(m.Value)) } if err := r.Close(); err != nil { log.Fatal("failed to close reader:", err) } ``` ## 日志记录 为了解 Reader/Writer 类型的操作情况,可以在创建时配置日志记录器。 ### Reader ``` func logf(msg string, a ...interface{}) { fmt.Printf(msg, a...) fmt.Println() } r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"}, Topic: "my-topic1", Partition: 0, Logger: kafka.LoggerFunc(logf), ErrorLogger: kafka.LoggerFunc(logf), }) ``` ### Writer ``` func logf(msg string, a ...interface{}) { fmt.Printf(msg, a...) fmt.Println() } w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092"), Topic: "topic", Logger: kafka.LoggerFunc(logf), ErrorLogger: kafka.LoggerFunc(logf), } ``` ## 测试 Kafka 后续版本中细微的行为变化导致了一些历史测试失败,如果你针对 Kafka 2.3.1 或更高版本运行,导出 `KAFKA_SKIP_NETTEST=1` 环境变量将跳过这些测试。 在 docker 中本地运行 Kafka ``` docker-compose up -d ``` 运行测试 ``` KAFKA_VERSION=2.3.1 \ KAFKA_SKIP_NETTEST=1 \ go test -race ./... ``` (或者)要清理缓存的测试结果并运行测试: ``` go clean -cache && make test ```
标签:EVTX分析, 日志审计