中文字幕av专区_日韩电影在线播放_精品国产精品久久一区免费式_av在线免费观看网站

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

如何使用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

彭州市| 昌邑市| 梁平县| 桦南县| 洪泽县| 吉林市| 百色市| 哈巴河县| 渝北区| 凉城县| 两当县| 绥滨县| 冷水江市| 肃北| 长治市| 东台市| 泰安市| 大田县| 台前县| 新干县| 蒙城县| 平原县| 辽源市| 松江区| 石林| 噶尔县| 宁乡县| 漠河县| 孟连| 宜良县| 繁昌县| 柳林县| 麻阳| 宜宾县| 敖汉旗| 宣威市| 饶阳县| 龙陵县| 和硕县| 桐城市| 花垣县|