最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Golang操作Kafka的實現(xiàn)示例

 更新時間:2023年02月19日 09:00:04   作者:YUHAOHAO  
本文主要介紹了Golang操作Kafka的實現(xiàn)示例,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

一.使用庫說明

Golang中連接kafka可以使用第三方庫:github.com/Shopify/sarama

二.Kafka Producer發(fā)送消息

package main?

import (
?? ?"fmt"
?? ?"github.com/Shopify/sarama"
)

func main() {
?? ?config := sarama.NewConfig()
?? ?config.Producer.RequiredAcks = sarama.WaitForAll // 發(fā)送完數(shù)據(jù)需要leader和follower都確認
?? ?config.Producer.Partitioner = sarama.NewRandomPartitioner ?//寫到隨機分區(qū)中,我們默認設置32個分區(qū)
?? ?config.Producer.Return.Successes = true // 成功交付的消息將在success channel返回

?? ?// 構造一個消息
?? ?msg := &sarama.ProducerMessage{}
?? ?msg.Topic = "task"
?? ?msg.Value = sarama.StringEncoder("producer kafka messages...")

?? ?// 連接kafka
?? ?client, err := sarama.NewSyncProducer([]string{"192.20.216.8:9092"}, config)
?? ?if err != nil {
?? ??? ?fmt.Println("Producer closed, err:", err)
?? ??? ?return
?? ?}
?? ?defer client.Close()

?? ?// 發(fā)送消息
?? ?pid, offset, err := client.SendMessage(msg)
?? ?if err != nil {
?? ??? ?fmt.Println("send msg failed, err:", err)
?? ??? ?return
?? ?}
?? ?fmt.Printf("pid:%v offset:%v\n", pid, offset)
}

三.Kafka Consumer消費消息

package main

import (
?? ?"fmt"
?? ?"github.com/Shopify/sarama"
?? ?"sync"
)

func main() {
?? ?var wg sync.WaitGroup
?? ?consumer, err := sarama.NewConsumer([]string{"192.20.216.8:9092"}, nil)
?? ?if err != nil {
?? ??? ?fmt.Println("Failed to start consumer: %s", err)
?? ??? ?return
?? ?}
?? ?partitionList, err := consumer.Partitions("task-status-data") // 通過topic獲取到所有的分區(qū)
?? ?if err != nil {
?? ??? ?fmt.Println("Failed to get the list of partition: ", err)
?? ??? ?return
?? ?}
?? ?fmt.Println(partitionList)

?? ?for partition := range partitionList{ // 遍歷所有的分區(qū)
?? ??? ?pc, err := consumer.ConsumePartition("task", int32(partition), sarama.OffsetNewest) // 針對每個分區(qū)創(chuàng)建一個分區(qū)消費者
?? ??? ?if err != nil {
?? ??? ??? ?fmt.Println("Failed to start consumer for partition %d: %s\n", partition, err)
?? ??? ?}
?? ??? ?wg.Add(1)
?? ??? ?go func(sarama.PartitionConsumer) { // 為每個分區(qū)開一個go協(xié)程取值
?? ??? ??? ?for msg := range pc.Messages() { // 阻塞直到有值發(fā)送過來,然后再繼續(xù)等待
?? ??? ??? ??? ?fmt.Printf("Partition:%d, Offset:%d, key:%s, value:%s\n", msg.Partition, msg.Offset, string(msg.Key), string(msg.Value))
?? ??? ??? ?}
?? ??? ??? ?defer pc.AsyncClose()
?? ??? ??? ?wg.Done()
?? ??? ?}(pc)
?? ?}
?? ?wg.Wait()
?? ?consumer.Close()
}

到此這篇關于Golang操作Kafka的實現(xiàn)示例的文章就介紹到這了,更多相關Golang操作Kafka內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • go-zero熔斷機制組件Breaker接口定義使用解析

    go-zero熔斷機制組件Breaker接口定義使用解析

    這篇文章主要為大家介紹了go-zero熔斷機制組件Breaker接口定義使用解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-05-05
  • golang time包下定時器的實現(xiàn)方法

    golang time包下定時器的實現(xiàn)方法

    定時器的實現(xiàn)大家應該都遇到過,最近在學習golang,所以下面這篇文章主要給大家介紹了關于golang time包下定時器的實現(xiàn)方法,文中通過示例代碼介紹的非常詳細,需要的朋友可以參考借鑒,下面來一起看看吧。
    2017-12-12
  • 深入理解Go設計模式之代理模式

    深入理解Go設計模式之代理模式

    代理模式是一種結構型設計模式,?其中代理控制著對于原對象的訪問,?并允許在將請求提交給原對象的前后進行一些處理,從而增強原對象的邏輯處理,這篇文章主要來學習一下代理模式的構成和用法,需要的朋友可以參考下
    2023-05-05
  • 詳解go-admin在線開發(fā)平臺學習(安裝、配置、啟動)

    詳解go-admin在線開發(fā)平臺學習(安裝、配置、啟動)

    這篇文章主要介紹了go-admin在線開發(fā)平臺學習(安裝、配置、啟動),本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-02-02
  • 詳解Go語言變量作用域

    詳解Go語言變量作用域

    這篇文章主要介紹了Go 語言變量作用域的相關資料,幫助大家更好的理解和學習使用go語言,感興趣的朋友可以了解下
    2021-03-03
  • GO中Json解析的幾種方式

    GO中Json解析的幾種方式

    本文主要介紹了GO中Json解析的幾種方式,詳細的介紹了幾種方法,?文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2024-01-01
  • 基于Go+OpenCV實現(xiàn)人臉識別功能的詳細示例

    基于Go+OpenCV實現(xiàn)人臉識別功能的詳細示例

    OpenCV是一個強大的計算機視覺庫,提供了豐富的圖像處理和計算機視覺算法,本文將向你介紹在Mac上安裝OpenCV的步驟,并演示如何使用Go的OpenCV綁定庫進行人臉識別,需要的朋友可以參考下
    2023-07-07
  • go編程中go-sql-driver的離奇bug解決記錄分析

    go編程中go-sql-driver的離奇bug解決記錄分析

    這篇文章主要為大家介紹了go編程中go-sql-driver的離奇bug解決記錄分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-05-05
  • Go語言基礎學習教程

    Go語言基礎學習教程

    這篇文章主要介紹了Go語言基礎知識,包括基本語法、語句、數(shù)組等的定義與用法,需要的朋友可以參考下
    2016-07-07
  • GoLang jwt無感刷新與SSO單點登錄限制解除方法詳解

    GoLang jwt無感刷新與SSO單點登錄限制解除方法詳解

    這篇文章主要介紹了GoLang jwt無感刷新與SSO單點登錄限制解除方法,JWT是一個簽名的JSON對象,通常用作Oauth2的Bearer token,JWT包括三個用.分割的部分。本文將利用JWT進行認證和加密,感興趣的可以了解一下
    2023-03-03

最新評論

师宗县| 武功县| 金川县| 垦利县| 乐昌市| 丹阳市| 元江| 吕梁市| 栾城县| 红桥区| 麻城市| 井陉县| 西宁市| 泰州市| 屏南县| 五指山市| 特克斯县| 焦作市| 额济纳旗| 永清县| 桦南县| 辽宁省| 彭水| 吉隆县| 武安市| 江口县| 绥芬河市| 通渭县| 吉隆县| 茶陵县| 湛江市| 保德县| 习水县| 陇西县| 威信县| 郎溪县| 宁河县| 通城县| 芷江| 安福县| 宝丰县|