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

go操作Kafka使用示例詳解

 更新時間:2022年12月05日 09:13:39   作者:qi66  
這篇文章主要為大家介紹了go操作Kafka使用示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

1. Kafka介紹

1.1 Kafka是什么

kafka使用scala開發(fā),支持多語言客戶端(c++、java、python、go等)

Kafka最先由LinkedIn公司開發(fā),之后成為Apache的頂級項目。

Kafka是一個分布式的、分區(qū)化、可復(fù)制提交的日志服務(wù)

LinkedIn使用Kafka實現(xiàn)了公司不同應(yīng)用程序之間的松耦和,那么作為一個可擴展、高可靠的消息系統(tǒng) 支持高Throughput的應(yīng)用

scale out:無需停機即可擴展機器

持久化:通過將數(shù)據(jù)持久化到硬盤以及replication防止數(shù)據(jù)丟失

支持online和offline的場景

1.2 Kafka的特點

Kafka是分布式的,其所有的構(gòu)件borker(服務(wù)端集群)、producer(消息生產(chǎn))、consumer(消息消費者)都可以是分布式的。

在消息的生產(chǎn)時可以使用一個標(biāo)識topic來區(qū)分,且可以進行分區(qū);每一個分區(qū)都是一個順序的、不可變的消息隊列, 并且可以持續(xù)的添加。

同時為發(fā)布和訂閱提供高吞吐量。據(jù)了解,Kafka每秒可以生產(chǎn)約25萬消息(50 MB),每秒處理55萬消息(110 MB)。

消息被處理的狀態(tài)是在consumer端維護,而不是由server端維護。當(dāng)失敗時能自動平衡

1.3 常用的場景

監(jiān)控:主機通過Kafka發(fā)送與系統(tǒng)和應(yīng)用程序健康相關(guān)的指標(biāo),然后這些信息會被收集和處理從而創(chuàng)建監(jiān)控儀表盤并發(fā)送警告。

消息隊列: 應(yīng)用程度使用Kafka作為傳統(tǒng)的消息系統(tǒng)實現(xiàn)標(biāo)準(zhǔn)的隊列和消息的發(fā)布—訂閱,例如搜索和內(nèi)容提要(Content Feed)。比起大多數(shù)的消息系統(tǒng)來說,Kafka有更好的吞吐量,內(nèi)置的分區(qū),冗余及容錯性,這讓Kafka成為了一個很好的大規(guī)模消息處理應(yīng)用的解決方案。消息系統(tǒng) 一般吞吐量相對較低,但是需要更小的端到端延時,并嘗嘗依賴于Kafka提供的強大的持久性保障。在這個領(lǐng)域,Kafka足以媲美傳統(tǒng)消息系統(tǒng),如ActiveMR或RabbitMQ

站點的用戶活動追蹤: 為了更好地理解用戶行為,改善用戶體驗,將用戶查看了哪個頁面、點擊了哪些內(nèi)容等信息發(fā)送到每個數(shù)據(jù)中心的Kafka集群上,并通過Hadoop進行分析、生成日常報告。

流處理:保存收集流數(shù)據(jù),以提供之后對接的Storm或其他流式計算框架進行處理。很多用戶會將那些從原始topic來的數(shù)據(jù)進行 階段性處理,匯總,擴充或者以其他的方式轉(zhuǎn)換到新的topic下再繼續(xù)后面的處理。例如一個文章推薦的處理流程,可能是先從RSS數(shù)據(jù)源中抓取文章的內(nèi) 容,然后將其丟入一個叫做“文章”的topic中;后續(xù)操作可能是需要對這個內(nèi)容進行清理,比如回復(fù)正常數(shù)據(jù)或者刪除重復(fù)數(shù)據(jù),最后再將內(nèi)容匹配的結(jié)果返 還給用戶。這就在一個獨立的topic之外,產(chǎn)生了一系列的實時數(shù)據(jù)處理的流程。

日志聚合:使用Kafka代替日志聚合(log aggregation)。日志聚合一般來說是從服務(wù)器上收集日志文件,然后放到一個集中的位置(文件服務(wù)器或HDFS)進行處理。然而Kafka忽略掉 文件的細節(jié),將其更清晰地抽象成一個個日志或事件的消息流。這就讓Kafka處理過程延遲更低,更容易支持多數(shù)據(jù)源和分布式數(shù)據(jù)處理。比起以日志為中心的 系統(tǒng)比如Scribe或者Flume來說,Kafka提供同樣高效的性能和因為復(fù)制導(dǎo)致的更高的耐用性保證,以及更低的端到端延遲

持久性日志:Kafka可以為一種外部的持久性日志的分布式系統(tǒng)提供服務(wù)。這種日志可以在節(jié)點間備份數(shù)據(jù),并為故障節(jié)點數(shù)據(jù)回復(fù)提供一種重新同步的機制。Kafka中日志壓縮功能為這種用法提供了條件。在這種用法中,Kafka類似于Apache BookKeeper項目。

1.4 Kafka中包含以下基礎(chǔ)概念

1.Topic(話題):Kafka中用于區(qū)分不同類別信息的類別名稱。由producer指定

2.Producer(生產(chǎn)者):將消息發(fā)布到Kafka特定的Topic的對象(過程)

3.Consumers(消費者):訂閱并處理特定的Topic中的消息的對象(過程)

4.Broker(Kafka服務(wù)集群):已發(fā)布的消息保存在一組服務(wù)器中,稱之為Kafka集群。集群中的每一個服務(wù)器都是一個代理(Broker). 消費者可以訂閱一個或多個話題,并從Broker拉數(shù)據(jù),從而消費這些已發(fā)布的消息。

5.Partition(分區(qū)):Topic物理上的分組,一個topic可以分為多個partition,每個partition是一個有序的隊列。partition中的每條消息都會被分配一個有序的id(offset)

Message:消息,是通信的基本單位,每個producer可以向一個topic(主題)發(fā)布一些消息。

1.5 消息

消息由一個固定大小的報頭和可變長度但不透明的字節(jié)陣列負載。報頭包含格式版本和CRC32效驗和以檢測損壞或截斷

1.6 消息格式

    1. 4 byte CRC32 of the message
    2. 1 byte "magic" identifier to allow format changes, value is 0 or 1
    3. 1 byte "attributes" identifier to allow annotations on the message independent of the version
       bit 0 ~ 2 : Compression codec
           0 : no compression
           1 : gzip
           2 : snappy
           3 : lz4
       bit 3 : Timestamp type
           0 : create time
           1 : log append time
       bit 4 ~ 7 : reserved
    4. (可選) 8 byte timestamp only if "magic" identifier is greater than 0
    5. 4 byte key length, containing length K
    6. K byte key
    7. 4 byte payload length, containing length V
    8. V byte payload

2. Kafka深層介紹

2.1 架構(gòu)介紹

Producer:Producer即生產(chǎn)者,消息的產(chǎn)生者,是消息的?口。

kafka cluster:kafka集群,一臺或多臺服務(wù)?組成

  • Broker:Broker是指部署了Kafka實例的服務(wù)?節(jié)點。每個服務(wù)?上有一個或多個kafka的實 例,我們姑且認(rèn)為每個broker對應(yīng)一臺服務(wù)?。每個kafka集群內(nèi)的broker都有一個不重復(fù)的 編號,如圖中的broker-0、broker-1等……
  • Topic:消息的主題,可以理解為消息的分類,kafka的數(shù)據(jù)就保存在topic。在每個broker上 都可以創(chuàng)建多個topic。實際應(yīng)用中通常是一個業(yè)務(wù)線建一個topic。
  • Partition:Topic的分區(qū),每個topic可以有多個分區(qū),分區(qū)的作用是做負載,提高kafka的吞 吐量。同一個topic在不同的分區(qū)的數(shù)據(jù)是不重復(fù)的,partition的表現(xiàn)形式就是一個一個的?件夾!
  • Replication:每一個分區(qū)都有多個副本,副本的作用是做備胎。當(dāng)主分區(qū)(Leader)故障的 時候會選擇一個備胎(Follower)上位,成為Leader。在kafka中默認(rèn)副本的最大數(shù)量是10 個,且副本的數(shù)量不能大于Broker的數(shù)量,follower和leader絕對是在不同的機器,同一機 ?對同一個分區(qū)也只可能存放一個副本(包括自己)。

Consumer:消費者,即消息的消費方,是消息的出口。

  • Consumer Group:我們可以將多個消費組組成一個消費者組,在kafka的設(shè)計中同一個分 區(qū)的數(shù)據(jù)只能被消費者組中的某一個消費者消費。同一個消費者組的消費者可以消費同一個 topic的不同分區(qū)的數(shù)據(jù),這也是為了提高kafka的吞吐量!

2.2 ?作流程

我們看上?的架構(gòu)圖中,producer就是生產(chǎn)者,是數(shù)據(jù)的入口。Producer在寫入數(shù)據(jù)的時候會把數(shù)據(jù) 寫入到leader中,不會直接將數(shù)據(jù)寫入follower!那leader怎么找呢?寫入的流程又是什么樣的呢?我 們看下圖:

1.?產(chǎn)者從Kafka集群獲取分區(qū)leader信息

2.?產(chǎn)者將消息發(fā)送給leader

3.leader將消息寫入本地磁盤

4.follower從leader拉取消息數(shù)據(jù)

5.follower將消息寫入本地磁盤后向leader發(fā)送ACK

6.leader收到所有的follower的ACK之后向生產(chǎn)者發(fā)送ACK

2.3 選擇partition的原則

那在kafka中,如果某個topic有多個partition,producer?怎么知道該將數(shù)據(jù)發(fā)往哪個partition呢? kafka中有幾個原則:

1.partition在寫入的時候可以指定需要寫入的partition,如果有指定,則寫入對應(yīng)的partition。

2.如果沒有指定partition,但是設(shè)置了數(shù)據(jù)的key,則會根據(jù)key的值hash出一個partition。

3.如果既沒指定partition,又沒有設(shè)置key,則會采用輪詢?式,即每次取一小段時間的數(shù)據(jù)寫入某partition,下一小段的時間寫入下一個partition

2.4 ACK應(yīng)答機制

producer在向kafka寫入消息的時候,可以設(shè)置參數(shù)來確定是否確認(rèn)kafka接收到數(shù)據(jù),這個參數(shù)可設(shè)置 的值為 0,1,all

  • 0代表producer往集群發(fā)送數(shù)據(jù)不需要等到集群的返回,不確保消息發(fā)送成功。安全性最低但是效 率最高。
  • 1代表producer往集群發(fā)送數(shù)據(jù)只要leader應(yīng)答就可以發(fā)送下一條,只確保leader發(fā)送成功。
  • all代表producer往集群發(fā)送數(shù)據(jù)需要所有的follower都完成從leader的同步才會發(fā)送下一條,確保 leader發(fā)送成功和所有的副本都完成備份。安全性最?高,但是效率最低。

最后要注意的是,如果往不存在的topic寫數(shù)據(jù),kafka會?動創(chuàng)建topic,partition和replication的數(shù)量 默認(rèn)配置都是1。

2.5 Topic和數(shù)據(jù)?志

topic 是同?類別的消息記錄(record)的集合。在Kafka中,?個主題通常有多個訂閱者。對于每個 主題,Kafka集群維護了?個分區(qū)數(shù)據(jù)?志?件結(jié)構(gòu)如下:

每個partition都是?個有序并且不可變的消息記錄集合。當(dāng)新的數(shù)據(jù)寫?時,就被追加到partition的末 尾。在每個partition中,每條消息都會被分配?個順序的唯?標(biāo)識,這個標(biāo)識被稱為offset,即偏移 量。注意,Kafka只保證在同?個partition內(nèi)部消息是有序的,在不同partition之間,并不能保證消息 有序。

Kafka可以配置?個保留期限,?來標(biāo)識?志會在Kafka集群內(nèi)保留多?時間。Kafka集群會保留在保留 期限內(nèi)所有被發(fā)布的消息,不管這些消息是否被消費過。?如保留期限設(shè)置為兩天,那么數(shù)據(jù)被發(fā)布到 Kafka集群的兩天以內(nèi),所有的這些數(shù)據(jù)都可以被消費。當(dāng)超過兩天,這些數(shù)據(jù)將會被清空,以便為后 續(xù)的數(shù)據(jù)騰出空間。由于Kafka會將數(shù)據(jù)進?持久化存儲(即寫?到硬盤上),所以保留的數(shù)據(jù)??可 以設(shè)置為?個?較?的值。

2.6 Partition結(jié)構(gòu)

Partition在服務(wù)器上的表現(xiàn)形式就是?個?個的?件夾,每個partition的?件夾下?會有多組segment ?件,每組segment?件?包含 .index ?件、 .log ?件、 .timeindex ?件三個?件,其中 .log ? 件就是實際存儲message的地?,? .index 和 .timeindex ?件為索引?件,?于檢索消息。

2.7 消費數(shù)據(jù)

多個消費者實例可以組成?個消費者組,并??個標(biāo)簽來標(biāo)識這個消費者組。?個消費者組中的不同消 費者實例可以運?在不同的進程甚?不同的服務(wù)器上。

如果所有的消費者實例都在同?個消費者組中,那么消息記錄會被很好的均衡的發(fā)送到每個消費者實 例。

如果所有的消費者實例都在不同的消費者組,那么每?條消息記錄會被?播到每?個消費者實例。

舉個例?,如上圖所示?個兩個節(jié)點的Kafka集群上擁有?個四個partition(P0-P3)的topic。有兩個 消費者組都在消費這個topic中的數(shù)據(jù),消費者組A有兩個消費者實例,消費者組B有四個消費者實例。 從圖中我們可以看到,在同?個消費者組中,每個消費者實例可以消費多個分區(qū),但是每個分區(qū)最多只 能被消費者組中的?個實例消費。也就是說,如果有?個4個分區(qū)的主題,那么消費者組中最多只能有4 個消費者實例去消費,多出來的都不會被分配到分區(qū)。其實這也很好理解,如果允許兩個消費者實例同 時消費同?個分區(qū),那么就?法記錄這個分區(qū)被這個消費者組消費的offset了。如果在消費者組中動態(tài) 的上線或下線消費者,那么Kafka集群會?動調(diào)整分區(qū)與消費者實例間的對應(yīng)關(guān)系。

3. 操作Kafka

3.1 sarama

Go語言中連接kafka使用第三方庫: github.com/Shopify/sarama。

3.2 下載及安裝

    go get github.com/Shopify/sarama

注意事項: sarama v1.20之后的版本加入了zstd壓縮算法,需要用到cgo,在Windows平臺編譯時會提示類似如下錯誤: github.com/DataDog/zstd exec: "gcc":executable file not found in %PATH% 所以在Windows平臺請使用v1.19版本的sarama。(如果不會版本控制請查看博客里面的go module章節(jié))

3.3 連接kafka發(fā)送消息

package main
import (
    "fmt"
    "github.com/Shopify/sarama"
)
// 基于sarama第三方庫開發(fā)的kafka client
func main() {
    config := sarama.NewConfig()
    config.Producer.RequiredAcks = sarama.WaitForAll          // 發(fā)送完數(shù)據(jù)需要leader和follow都確認(rèn)
    config.Producer.Partitioner = sarama.NewRandomPartitioner // 新選出一個partition
    config.Producer.Return.Successes = true                   // 成功交付的消息將在success channel返回
    // 構(gòu)造一個消息
    msg := &sarama.ProducerMessage{}
    msg.Topic = "web_log"
    msg.Value = sarama.StringEncoder("this is a test log")
    // 連接kafka
    client, err := sarama.NewSyncProducer([]string{"127.0.0.1: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)
}

3.4 連接kafka消費消息

package main
import (
    "fmt"
    "github.com/Shopify/sarama"
)
// kafka consumer
func main() {
    consumer, err := sarama.NewConsumer([]string{"127.0.0.1:9092"}, nil)
    if err != nil {
        fmt.Printf("fail to start consumer, err:%v\n", err)
        return
    }
    partitionList, err := consumer.Partitions("web_log") // 根據(jù)topic取到所有的分區(qū)
    if err != nil {
        fmt.Printf("fail to get list of partition:err%v\n", err)
        return
    }
    fmt.Println(partitionList)
    for partition := range partitionList { // 遍歷所有的分區(qū)
        // 針對每個分區(qū)創(chuàng)建一個對應(yīng)的分區(qū)消費者
        pc, err := consumer.ConsumePartition("web_log", int32(partition), sarama.OffsetNewest)
        if err != nil {
            fmt.Printf("failed to start consumer for partition %d,err:%v\n", partition, err)
            return
        }
        defer pc.AsyncClose()
        // 異步從每個分區(qū)消費信息
        go func(sarama.PartitionConsumer) {
            for msg := range pc.Messages() {
                fmt.Printf("Partition:%d Offset:%d Key:%v Value:%v", msg.Partition, msg.Offset, msg.Key, msg.Value)
            }
        }(pc)
    }
}

以上就是go操作Kfaka使用示例詳解的詳細內(nèi)容,更多關(guān)于go操作Kfaka的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 關(guān)于Golang變量初始化/類型推斷/短聲明的問題

    關(guān)于Golang變量初始化/類型推斷/短聲明的問題

    這篇文章主要介紹了關(guān)于Golang變量初始化/類型推斷/短聲明的問題,本文給大家介紹的非常詳細,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-02-02
  • Go高級特性之并發(fā)處理http詳解

    Go高級特性之并發(fā)處理http詳解

    Golang?作為一種高效的編程語言,提供了多種方法來實現(xiàn)并發(fā)發(fā)送?HTTP?請求,本文將深入探討?Golang?中并發(fā)發(fā)送?HTTP?請求的最佳技術(shù)和實踐,希望對大家有所幫助
    2024-02-02
  • golang如何解決go get命令無響應(yīng)問題

    golang如何解決go get命令無響應(yīng)問題

    文章介紹了在Go語言中處理由于官方庫被封禁導(dǎo)致依賴下載失敗的方法,包括設(shè)置代理和直接克隆依賴包到GOPATH/src下
    2024-12-12
  • go語言的sql包原理與用法分析

    go語言的sql包原理與用法分析

    這篇文章主要介紹了go語言的sql包原理與用法,較為詳細的分析了Go語言里sql包的結(jié)構(gòu)、相關(guān)函數(shù)與使用方法,需要的朋友可以參考下
    2016-07-07
  • Go 面向包新提案透明文件夾必要性分析

    Go 面向包新提案透明文件夾必要性分析

    這篇文章主要為大家介紹了Go 面向包新提案,透明文件夾必要性分析,看看是否合適加進 Go 特性中,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-11-11
  • 深入解析Go語言中HTTP請求處理的底層實現(xiàn)

    深入解析Go語言中HTTP請求處理的底層實現(xiàn)

    本文將詳細介紹?Go?語言中?HTTP?請求處理的底層機制,包括工作流程、創(chuàng)建?Listen?Socket?監(jiān)聽端口、接收客戶端請求并建立連接以及處理客戶端請求并返回響應(yīng)等,需要的朋友可以參考下
    2023-05-05
  • Go語言內(nèi)建函數(shù)cap的實現(xiàn)示例

    Go語言內(nèi)建函數(shù)cap的實現(xiàn)示例

    cap 是一個常用的內(nèi)建函數(shù),它用于獲取某些數(shù)據(jù)結(jié)構(gòu)的容量,本文主要介紹了Go語言內(nèi)建函數(shù)cap的實現(xiàn)示例,具有一定的參考價值,感興趣的可以了解一下
    2024-08-08
  • Golang使用http協(xié)議實現(xiàn)心跳檢測程序過程詳解

    Golang使用http協(xié)議實現(xiàn)心跳檢測程序過程詳解

    這篇文章主要介紹了Golang使用http協(xié)議實現(xiàn)心跳檢測程序過程,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧
    2023-03-03
  • Go經(jīng)典面試題匯總(填空+判斷)

    Go經(jīng)典面試題匯總(填空+判斷)

    這篇文章主要介紹了Go經(jīng)典面試題匯總(填空+判斷),本文章內(nèi)容詳細,具有很好的參考價值,希望對大家有所幫助,需要的朋友可以參考下
    2023-01-01
  • Golang實現(xiàn)JWT身份驗證的示例詳解

    Golang實現(xiàn)JWT身份驗證的示例詳解

    JWT(JSON Web Token)是一種開放標(biāo)準(zhǔn)(RFC 7519),用于在網(wǎng)絡(luò)應(yīng)用間安全地傳輸聲明,本文主要為大家詳細介紹了Golang實現(xiàn)JWT身份驗證的相關(guān)方法,希望對大家有所幫助
    2024-03-03

最新評論

麟游县| 荥阳市| 周至县| 报价| 宁蒗| 东明县| 新巴尔虎左旗| 潢川县| 衡山县| 共和县| 浮山县| 哈密市| 克东县| 临城县| 吴江市| 邵阳县| 拉孜县| 大宁县| 华安县| 城市| 饶河县| 彭阳县| 和顺县| 灯塔市| 灵石县| 仲巴县| 望江县| 潜江市| 托克逊县| 太谷县| 盈江县| 五台县| 乐业县| 武定县| 武山县| 萍乡市| 东台市| 秭归县| 开化县| 隆尧县| 无棣县|