首页
学习
活动
专区
工具
TVP
发布
社区首页 >专栏 >Golang 语言中 kafka 客户端库 sarama

Golang 语言中 kafka 客户端库 sarama

作者头像
frank.
发布2021-05-13 17:01:29
6.2K0
发布2021-05-13 17:01:29
举报

介绍

Apache Kafka 是一款开源的消息引擎系统。它在项目中的作用主要是削峰填谷和解耦。本文我们只介绍 Apache Kafka 的 Golang 客户端库 Sarama。Sarama 是 MIT 许可的 Apache Kafka 0.8 及更高版本的 Golang 客户端库。

如果读者朋友对 Apache Kafka 服务端还不了解,建议先阅读官方文档中的入门部分,本文使用的版本是 Apache Kafka 2.8。

02

生产者

我们可以使用 Sarama 库的 AsyncProducer 或 SyncProducer 生产消息。在大多数情况下首选使用 AsyncProducer 生产消息。它通过一个 channel 接收消息,并在后台尽可能高效的异步生产消息。

SyncProducer 发送 Kafka 消息后阻塞,直到接收到 ACK 确认。SyncProducer 有两个警告:它通常效率较低,并且实际的耐用性保证取决于 Producer.RequiredAcks 的配置值。在某些配置中,有时仍会丢失由 SyncProducer 确认的消息,但是使用比较简单。

为了读者朋友们容易理解,本文我们介绍 SyncProducer 作为生产者的使用方式。如果读者朋友想了解 AsyncProducer 作为生产者的使用方式,请参考官方文档。

使用 SyncProducer 作为生产者的示例代码:

func sendMessage (brokerAddr []string, config *sarama.Config, topic string, value sarama.Encoder) {
 producer, err := sarama.NewSyncProducer(brokerAddr, config)
 if err != nil {
  fmt.Println(err)
  return
 }
 defer func() {
  if err = producer.Close(); err != nil {
   fmt.Println(err)
   return
  }
 }()
 msg := &sarama.ProducerMessage{
  Topic: topic,
  Value: value,
 }
 partition, offset, err := producer.SendMessage(msg)
 if err != nil {
  fmt.Println(err)
  return
 }
 fmt.Printf("partition:%d offset:%d\n", partition, offset)
}

阅读上面这段代码,我们调用 NewSyncProducer() 创建一个新的 SyncProducer,给定 broker 地址和配置信息。调用 SendMessage() 生产给定的消息,并且仅在生产成功或失败时返回。它将返回分区(Partition)和生产的消息的偏移量(Offset),如果消息生产失败,则返回错误。

需要注意的是,为了避免泄露,必须在生产者上调用 Close(),因为当它超出范围时,可能不会自动垃圾回收。

03

消费者

我们可以使用 Sarama 库的消费者 Consumer 或消费者组 ConsumerGroup API 消费消息。为了读者朋友们容易理解,本文我们介绍使用 Consumer 消费消息。

Consumer 管理 PartitionConsumers,该 PartitionConsumers 处理来自 brokers 的 Kafka 消息。

Consumer 消费消息的示例代码:

func consumer (brokenAddr []string, topic string, partition int32, offset int64) {
 consumer, err := sarama.NewConsumer(brokenAddr, nil)
 if err != nil {
  fmt.Println(err)
  return
 }
 defer func() {
  if err = consumer.Close(); err != nil {
   fmt.Println(err)
   return
  }
 }()
 partitionConsumer, err := consumer.ConsumePartition(topic, partition, offset)
 if err != nil {
  fmt.Println(err)
  return
 }
 defer func() {
  if err = partitionConsumer.Close(); err != nil {
   fmt.Println(err)
   return
  }
 }()
 for msg := range partitionConsumer.Messages() {
  fmt.Printf("partition:%d offset:%d key:%s val:%s\n", msg.Partition, msg.Offset, msg.Key, msg.Value)
 }
}

阅读上面这段代码,我们调用 NewConsumer() 创建一个新的 consumer,给定 broker 地址和配置信息。调用 ConsumePartition() 创建 PartitionConsumer,给定 topic、partition 和 offset。PartitionConsumer 处理来自给定 topic 和 partition 的 Kafka 消息。

需要注意的是,为了防止泄露,必须调用 consumer 和 partitionConsumer 的 Close(),因为当它超出范围时,可能不会自动垃圾回收。

04

总结

本文主要介绍如何使用 Apache Kafka 的 Golang 语言客户端库 Sarama 生产和消费 Kafka 消息。关于生产者和消费者,分别列举了一个简单示例。除此之外,Sarama 库还提供了很多其它 Api,感兴趣的读者朋友可以阅读官方文档了解更多。

推荐阅读:

Go 使用标准库 net/http 包构建服务器

Go 使用标准库 sql 包和三方数据库驱动包操作 MySQL

Golang 语言三方库 lumberjack 日志切割组件怎么使用?

Golang 语言的值验证库 Validator 怎么使用?

Go team 开源项目 Go Cloud 使用的依赖注入工具 Wire 怎么使用?

参考资料: https://github.com/Shopify/sarama https://shopify.github.io/sarama/ https://pkg.go.dev/github.com/Shopify/sarama https://github.com/Shopify/sarama/wiki/Frequently-Asked-Questions https://kafka.apache.org/documentation/

本文参与 腾讯云自媒体分享计划,分享自微信公众号。
原始发表:2021-05-06,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 Go语言开发栈 微信公众号,前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体分享计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
相关产品与服务
云数据库 MySQL
腾讯云数据库 MySQL(TencentDB for MySQL)为用户提供安全可靠,性能卓越、易于维护的企业级云数据库服务。其具备6大企业级特性,包括企业级定制内核、企业级高可用、企业级高可靠、企业级安全、企业级扩展以及企业级智能运维。通过使用腾讯云数据库 MySQL,可实现分钟级别的数据库部署、弹性扩展以及全自动化的运维管理,不仅经济实惠,而且稳定可靠,易于运维。
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档