segmentio/kafka-go
GitHub: segmentio/kafka-go
一个纯 Go 实现的 Kafka 客户端库,通过贴合 Go 标准库理念的高底层 API 解决了现有 Go Kafka 客户端难用、依赖重或不支持 context 的问题。
Stars: 8592 | Forks: 853
# kafka-go [](https://circleci.com/gh/segmentio/kafka-go) [](https://goreportcard.com/report/github.com/segmentio/kafka-go) [](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 或更高版本。
## 连接 [](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 [](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 [](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分析, 日志审计