kafka-go:Go语言实现的Kafka客户端,支持高低级API与Go标准库接口

Kafka library in Go

分支48Tags69
文件最后提交记录最后更新时间
1 年前
4 年前
3 年前
2 年前
2 年前
4 年前
3 年前
5 年前
2 年前
3 年前
1 年前
5 年前
5 年前
4 年前
5 年前
4 年前
4 年前
4 年前
4 年前
4 年前
9 年前
2 年前
2 年前
4 年前
1 年前
4 年前
4 年前
4 年前
5 年前
3 年前
3 年前
4 年前
4 年前
2 年前
2 年前
3 年前
3 年前
5 年前
5 年前
2 年前
2 年前
3 年前
4 年前
6 年前
4 年前
4 年前
4 年前
8 年前
8 年前
4 年前
3 年前
3 年前
2 年前
3 年前
6 年前
5 年前
2 年前
2 年前
5 年前
4 年前
3 年前
3 年前
2 年前
2 年前
3 年前
3 年前
3 年前
5 年前
3 年前
3 年前
3 年前
3 年前
4 年前
4 年前
4 年前
3 年前
3 年前
4 年前
3 年前
4 年前
4 年前
2 年前
5 年前
5 年前
4 年前
1 年前
1 年前
4 年前
4 年前
5 年前
3 年前
4 年前
4 年前
4 年前
2 年前
2 年前
4 年前
4 年前
3 年前
3 年前
5 年前
5 年前
4 年前
3 年前
3 年前
3 年前
3 年前
4 年前
4 年前
4 年前
5 年前
4 年前
5 年前
5 年前
2 年前
2 年前
4 年前
3 年前
3 年前
4 年前
3 年前
5 年前
4 年前
4 年前
3 年前
3 年前
3 年前
3 年前
3 年前
4 年前
3 年前
3 年前
2 年前
2 年前
4 年前
4 年前
2 年前
1 年前
4 年前
6 年前
4 年前
2 年前
2 年前
6 年前
6 年前
4 年前
6 年前
4 年前
3 年前
3 年前
3 年前
5 年前
3 年前
4 年前
4 年前
1 年前
6 年前
4 年前
2 年前
2 年前

kafka-go

CircleCI Go Report Card GoDoc

动机

Segment 在日常运作中高度依赖于 Go 和 Kafka。然而,在撰写本文时,可用的 Go 客户端库对于 Kafka 的支持并不令人满意。可选方案如下:

  • sarama:这是最流行的选项,但使用起来颇为复杂。文档不足,API 暴露了 Kafka 协议的底层概念,并且不支持诸如上下文(context)这样的现代 Go 特性。此外,它采用指针传递所有值的方式导致了大量的动态内存分配,频繁的垃圾回收和较高的内存消耗。

  • confluent-kafka-go:这是一个基于 cgo 缠绕 librdkafka 的封装,意味着在所有使用此包的 Go 代码中引入了一个对 C 库的依赖。相较于 sarama,它的文档更加详尽,但仍缺乏对 Go 上下文的支持。

  • goka:这是一个较新的 Kafka 的 Go 客户端,专注于特定的使用模式。它提供了用于服务间消息总线而非有序事件日志的抽象接口,但这并非我们在 Segment 使用 Kafka 的常见场景。该包还依赖于 sarama 进行所有与 Kafka 的交互。

正是在这种背景下,kafka-go 应运而生。它提供了一整套从低到高的 API 来与 Kafka 进行交互,模仿并实现了 Go 标准库中的概念和接口,便于集成现有软件,易于上手。

注意:

为了更好地与我们新采纳的行为准则保持一致,kafka-go 项目已将其默认分支重命名为 main。关于我们的行为准则详情,请参阅 这份 文件。

Kafka 版本

kafka-go 目前针对 Kafka 0.10.1.0 至 2.7.1 版本进行了测试。虽然它也应与更高版本兼容,但是 Kafka API 中新增的功能可能尚未在客户端实现。

Go 版本

kafka-go 需要使用 Go 1.15 或更高版本。

连接 GoDoc

Conn 类型是 kafka-go 包的核心部分,它围绕原生网络连接,暴露一个可以访问 Kafka 服务器的低级别 API。

以下是一些展示如何典型使用连接对象的示例:

// 发送消息至主题
topic := "my-topic"
partition := 0

conn, err := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", topic, partition)
if err != nil {
    log.Fatal("无法建立与领导者连接:", err)
}

conn.SetWriteDeadline(time.Now().Add(10 * time.Second))
_, err = conn.WriteMessages(
    kafka.Message{Value: []byte("一条!")},
    kafka.Message{Value: []byte("两条!")},
    kafka.Message{Value: []byte("三条!")},
)
if err != nil {
    log.Fatal("写入消息失败:", err)
}

if err := conn.Close(); err != nil {
    log.Fatal("关闭写操作出错:", err)
}
// 接收消息自主题
topic := "my-topic"
partition := 0

conn, err := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", topic, partition)
if err != nil {
    log.Fatal("无法建立与领导者连接:", err)
}

conn.SetReadDeadline(time.Now().Add(10 * time.Second))
batch := conn.ReadBatch(10e3, 1e6) // 获取至少10KB,最多1MB的数据

b := make([]byte, 10e3) // 每条消息最大10KB
for {
    n, err := batch.Read(b)
    if err != nil {
        break
    }
    fmt.Println(string(b[:n]))
}

if err := batch.Close(); err != nil {
    log.Fatal("批量读取关闭错误:", err)
}

if err := conn.Close(); err != nil {
    log.Fatal("连接关闭失败:", err)
}

创建主题

默认情况下,Kafka 设定 auto.create.topics.enable='true' (在 bitnami/kafka Kafka Docker 映像中为 KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE='true')。如果这个值设为 'true' ,那么通过 kafka.DialLeader 可以作为副效果创建主题,例如:

// 当 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',则需要明确地创建主题,过程如下:

// 当 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())
}

经由非领导节点连接至领导节点

// 不直接使用 DialLeader 而是经由现有的非领导者节点连接至 Kafka 领导者
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()

列举主题

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)的理想基石。

读取器 GoDoc

Readerkafka-go 包提供的另一个概念,旨在简化单个主题-分区组合消费的常见用例。 Reader 同时自动处理重新连接和偏移管理问题,并提供支持 Go 上下文中异步取消和超时机制的 API。

重要的是在进程退出时调用 Close() 方法来关闭 Reader 。Kafka 服务器需要一个优雅断开连接的过程来阻止其继续尝试向已连接的客户端发送消息。上述示例不会在进程被 SIGINT(shell 下的 Ctrl-C)或 SIGTERM(如 Docker 停止命令或 Kubernetes 重启动作)终止时调用 Close() 。这可能会导致当同一个主题的新读取器连接时出现延迟(比如启动新进程或运行新容器)。可以通过设置 signal.Notify 处理器实现在进程关机时关闭读取器。

// 创建一个新的读取器,从主题 A,分区 0,偏移量 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("偏移量 %d 的消息: %s = %s\n", m.Offset, string(m.Key), string(m.Value))
}

if err := r.Close(); err != nil {
    log.Fatal("读取器关闭失败:", err)
}

消费者群组

kafka-go也支持Kafka的消费者群组和由代理管理的偏移量。 要启用消费者群组,只需在ReaderConfig中指定GroupID。

使用消费者群组时,ReadMessage会自动提交偏移量。

// 创建一个消费"topic-A"的新阅读器
r := kafka.NewReader(kafka.ReaderConfig{
   Brokers:   {"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("消息位于 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("关闭阅读器失败:", err)
}

使用消费者群组时有一些限制:

  • 设置GroupID后调用(*Reader).SetOffset会返回错误
  • 当GroupID设置后,(*Reader).Offset总是返回-1
  • 当GroupID设置后,(*Reader).Lag总是返回-1
  • 当GroupID设置后,(*Reader).ReadLag会返回错误
  • 当GroupID设置后,(*Reader).Stats的分区总是-1

显式提交

kafka-go还支持显式提交。而不是调用ReadMessage,先调用FetchMessage,然后调用CommitMessages

ctx := context.Background()
for {
    m, err := r.FetchMessage(ctx)
    if err != nil {
        break
    }
    fmt.Printf("消息位于 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("提交消息失败:", err)
    }
}

在消费者群组中提交消息时,给定主题/分区的最高偏移量的消息决定了该分区提交的偏移量值。例如,如果通过调用FetchMessage获取了单个分区的偏移量为1、2和3的消息,则调用CommitMessages时传入偏移量为3也会提交该分区上偏移量为1和2的消息。

管理提交

默认情况下,CommitMessages会同步地将偏移量提交到Kafka。为了提高性能,可以通过在ReaderConfig中设置CommitInterval来定期向Kafka提交偏移量。

// 创建一个消费"topic-A"的新阅读器
r := kafka.NewReader(kafka.ReaderConfig{
    Brokers:        {"localhost:9092", "localhost:9093", "localhost:9094"},
    GroupID:        "consumer-group-id",
    Topic:          "topic-A",
    MaxBytes:       10e6, // 10MB
    CommitInterval: time.Second, // 每秒刷新一次提交
})

写入器 GoDoc

要向Kafka生产消息,程序可以使用低级的ConnAPI,但包还提供了一个更高级别的Writer类型,这在大多数情况下更适合使用,因为它提供了额外的功能:

  • 错误发生时的自动重试和重新连接。
  • 可配置的消息在可用分区间的分布。
  • 同步或异步写入消息到Kafka。
  • 使用上下文进行异步取消。
  • 关闭时刷新待处理消息以支持平滑关闭。
  • 在发布消息之前创建缺失的主题。 *注意!*这是在版本v0.4.30之前的默认行为。
// 创建一个写入器,它向"topic-A"发送消息,使用最少字节分布
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("你好世界!"),
    },
    kafka.Message{
        Key:   []byte("Key-B"),
        Value: []byte("一!"),
    },
    kafka.Message{
        Key:   []byte("Key-C"),
        Value: []byte("二!"),
    },
)
if err != nil {
    log.Fatal("写入消息失败:", err)
}

if err := w.Close(); err != nil {
    log.Fatal("关闭写入器失败:", err)
}

发布前创建缺失的主题

// 创建一个写入器,将其发送到"topic-A"。
// 如果主题不存在,将会创建。
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("你好世界!"),
    },
    {
        Key:   []byte("Key-B"),
        Value: []byte("一!"),
    },
    {
        Key:   []byte("Key-C"),
        Value: []byte("二!"),
    },
}

var err error
const retries = 3
for i := 0; i < retries; i++ {
    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()
    
    // 在发送消息前尝试创建主题
    err = w.WriteMessages(ctx, messages...)
    if errors.Is(err, kafka.LeaderNotAvailable) || errors.Is(err, context.DeadlineExceeded) {
        time.Sleep(250 * time.Millisecond)
        continue
    }

    if err != nil {
        log.Fatalf("意外的错误 %v", err)
    }
    break
}

if err := w.Close(); err != nil {
    log.Fatal("关闭写入器失败:", err)
}

向多个主题写入

通常,WriterConfig.Topic用于初始化单主题写入器。如果不设置这个配置,你可以通过在Message.Topic中为每个消息定义主题。

w := &kafka.Writer{
    Addr:     kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"),
    // 注意:这里不定义Topic,每个Message需自己定义。
    Balancer: &kafka.LeastBytes{},
}

err := w.WriteMessages(context.Background(),
    // 注意:每个Message都有Topic,否则会出错。
    kafka.Message{
        Topic: "topic-A",
        Key:   []byte("Key-A"),
        Value: []byte("你好世界!"),
    },
    kafka.Message{
        Topic: "topic-B",
        Key:   []byte("Key-B"),
        Value: []byte("一!"),
    },
    kafka.Message{
        Topic: "topic-C",
        Key:   []byte("Key-C"),
        Value: []byte("二!"),
    },
)
if err != nil {
    log.Fatal("写入消息失败:", err)
}

if err := w.Close(); err != nil {
    log.Fatal("关闭写入器失败:", err)
}

注意: 这两种模式是互斥的,如果你设置了Writer.Topic,就不应在要写入的消息上明确定义Message.Topic。相反,当没有为写入器定义主题时,也不应如此。如果检测到这种模糊性,Writer将返回错误。

兼容其他客户端

Sarama

如果你正从Sarama切换,并需要/希望使用相同的分区算法,你可以使用kafka.Hash平衡器或kafka.ReferenceHash平衡器:

  • kafka.Hash = sarama.NewHashPartitioner
  • kafka.ReferenceHash = sarama.NewReferenceHashPartitioner

kafka.Hashkafka.ReferenceHash 平衡器会将消息路由到与上述两个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平衡器可以获得与librdkafka的默认consistent_random分区策略相同的行为。

w := &kafka.Writer{
    Addr:     kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"),
    Topic:    "topic-A",
    Balancer: kafka.CRC32Balancer{},
}

Java

使用kafka.Murmur2Balancer平衡器可以获得与官方Java客户端的默认分发器相同的行为。请注意,Java类允许直接指定分区,而本库不允许这样做。

w := &kafka.Writer{
    Addr:     kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"),
    Topic:    "topic-A",
    Balancer: kafka.Murmur2Balancer{},
}

压缩功能

可对Writer启用压缩功能,只需设置Compression字段:

w := &kafka.Writer{
	Addr:        kafka.TCP("localhost:9092", "localhost:9093", "localhost:9094"),
	Topic:       "topic-A",
	Compression: kafka.Snappy,
}

Reader会自动识别接收到的消息是否被压缩,通过检查消息属性来确定。然而,所有预期编解码器对应的包都需导入以确保它们正确加载。

注意:在0.4版本前,程序需要导入压缩包来安装编解码器并支持从Kafka读取压缩消息。现在不再需要这样做,压缩包的导入已成为无操作过程。

TLS 支持

无论是在最基本的连接类型中还是在Reader/Writer配置里,你都可以指定TLS支持的拨号选项。如果TLS字段为空,则不会使用TLS建立连接。 重要提示: 如果未在连接/阅读器/写入器上配置TLS而尝试接入启用了TLS的Kafka集群,可能会导致模糊不清的io.ErrUnexpectedEOF错误。

连接示例

dialer := &kafka.Dialer{
    Timeout:   10 * time.Second,
    DualStack: true,
    TLS:       &tls.Config{...tls 配置...},
}

conn, err := dialer.DialContext(ctx, "tcp", "localhost:9093")

读者端(Reader)

dialer := &kafka.Dialer{
    Timeout:   10 * time.Second,
    DualStack: true,
    TLS:       &tls.Config{...tls 配置...},
}

r := kafka.NewReader(kafka.ReaderConfig{
    Brokers:        []string{"localhost:9092", "localhost:9093", "localhost:9094"},
    GroupID:        "consumer-group-id",
    Topic:          "topic-A",
    Dialer:         dialer,
})

写入者端(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 配置...},
}

w := kafka.NewWriter(kafka.WriterConfig{
	Brokers: []string{"localhost:9092", "localhost:9093", "localhost:9094"},
	Topic:   "topic-A",
	Balancer: &kafka.Hash{},
	Dialer:   dialer,
})

请注意,kafka.NewWriterkafka.WriterConfig已被标记为废弃,在未来的版本中将会移除。

SASL 支持

可以在Dialer中指定一个选项来使用SASL认证。可以直接使用Dialer打开Conn,或者将其作为参数传递给ReaderWriter各自所在的配置对象中。如果SASLMechanism字段为空,则不会执行SASL认证。

SASL 认证类型

Plain

mechanism := plain.Mechanism{
    Username: "username",
    Password: "password",
}

SCRAM

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)
}

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,
}

客户端(Client)

mechanism, err := scram.Mechanism(scram.SHA512, "username", "password")
if err != nil {
    panic(err)
}

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: 处理消息
    fmt.Printf("offset %d 上的消息: %s = %s\n", m.Offset, string(m.Key), string(m.Value))
}

if err := r.Close(); err != nil {
    log.Fatal("无法关闭读取器:", err)
}

日志记录

为了洞察ReaderWriter类型的操作详情,可在创建时为其配置日志记录器。

读者端(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