温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

如何使用Golang语言中的kafka和Sarama

发布时间:2021-09-13 14:50:52 来源:亿速云 阅读:195 作者:柒染 栏目:web开发

这篇文章给大家介绍如何使用Golang语言中的kafka和Sarama,内容非常详细,感兴趣的小伙伴们可以参考借鉴,希望对大家能有所帮助。

01、介绍

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

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

如何使用Golang语言中的kafka和Sarama

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(),因为当它超出范围时,可能不会自动垃圾回收。

关于如何使用Golang语言中的kafka和Sarama就分享到这里了,希望以上内容可以对大家有一定的帮助,可以学到更多知识。如果觉得文章不错,可以把它分享出去让更多的人看到。

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI