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

Golang中Kafka的重復消費和消息丟失問題的解決方案

 更新時間:2023年08月29日 09:46:51   作者:童話ing  
在Kafka中無論是生產(chǎn)者發(fā)送消息到Kafka集群還是消費者從Kafka集群中拉取消息,都是容易出現(xiàn)問題的,比較典型的就是消費端的重復消費問題、生產(chǎn)端和消費端產(chǎn)生的消息丟失問題,下面將對這兩個問題出現(xiàn)的場景以及常見的解決方案進行講解

前言

在Kafka中,生產(chǎn)者(Producer)和消費者(Consumer)是通過發(fā)布訂閱模式進行協(xié)作的,生產(chǎn)者將消息發(fā)送到Kafka集群,而消費者從Kafka集群中拉取消息進行消費,無論是生產(chǎn)者發(fā)送消息到Kafka集群還是消費者從Kafka集群中拉取消息進行消費,都是容易出現(xiàn)問題的,比較典型的就是消費端的重復消費問題、生產(chǎn)端和消費端產(chǎn)生的消息丟失問題。下面將對這兩個問題出現(xiàn)的場景以及常見的解決方案進行講解。

一、重復消費

1.1 重復消費出現(xiàn)的場景

重復消費出現(xiàn)的常見場景主要分為兩種,一種是 Consumer在消費過程中,應用進程被強制kill掉或者發(fā)生異常退出(掛掉…),另一種則是Consumer消費的時間過長。

1.1.1 Consumer消費過程中,進程掛掉/異常退出

在Kafka消費端的使用中,位移(Offset)的提交有兩種方式,自動提交和手動提交。自動提交情況下,當消費者拉取一批消息進行消費后,需要進行Offset的提交,在消費端提交Offset之前,Consumer掛掉了,當Consumer重啟后再次拉取Offset,這時候拉取的依然是掛掉之前消費的Offset,因此造成重復消費的問題。在手動提交模式下,在提交代碼調(diào)用之前,Consumer掛掉也會造成重復消費。

1.1.2 消費者消費時間過長

Kafka消費端的參數(shù)max.poll.interval.ms定義了兩次poll的最大間隔,它的默認值是 5 分鐘,表示 Consumer 如果在 5 分鐘之內(nèi)無法消費完 poll方法返回的消息,那么Consumer 會主動發(fā)起“離開組”的請求。

在離開消費組后,開始Rebalance,因此提交Offset失敗。之后重新Rebalance,消費者再次分配Partition后,再次poll拉取消息依然從之前消費過的消息處開始消費,這樣就造成重復消費。而且若不解決消費單次消費時間過長的問題,這部分消息可能會一直重復消費。

整體上來說,如果我們在消費中將消息數(shù)據(jù)處理入庫,但是在執(zhí)行Offset提交時,Kafka宕機或者網(wǎng)絡原因等無法提交Offset,當我們重啟服務或者Rebalance過程觸發(fā),Consumer將再次消費此消息數(shù)據(jù)。

1.2 重復消費解決方案

1.2.1 針對于消費端掛掉等原因造成的重復消費問題

這部分主要集中在消費端的編碼層面,需要我們在設計代碼時以冪等性的角度進行開發(fā)設計,保證同一數(shù)據(jù)無論進行多少次消費,所造成的結(jié)果都一樣。處理方式可以在消息體中添加唯一標識(比如將消息生成md5保存到Mysql或者是Redis中,在處理消息前先檢查下Mysql/Redis是否已經(jīng)處理過該消息了),消費端進行確認此唯一標識是否已經(jīng)消費過,如果消費過,則不進行之后處理。從而盡可能的避免了重復消費。
冪等角度大概兩種實現(xiàn):

  • 將唯一標識存入第三方介質(zhì)(如Redis),要操作數(shù)據(jù)的時候先判斷第三方介質(zhì)(數(shù)據(jù)庫或者緩存)有沒有這個唯一標識。
  • 將版本號(offset)存入到數(shù)據(jù)里面,然后再要操作數(shù)據(jù)的時候用這個版本號做樂觀鎖,當版本號大于原先的才能操作。

1.2.2 針對于Consumer消費時間過長帶來的重復消費問題

  • 提高單條消息的處理速度。例如對消息處理中比較耗時的操作可通過異步的方式進行處理、利用多線程處理等。
  • 其次,在縮短單條消息消費時常的同時,根據(jù)實際場景可將max.poll.interval.ms值設置大一點,避免不必要的rebalance,此外可適當減小max.poll.records的值,默認值是500,可根據(jù)實際消息速率適當調(diào)小。

二、消息丟失

在Kafka中,消息丟失在Kafka的生產(chǎn)端和消費端都會出現(xiàn)。在此之前我們先來了解一下生產(chǎn)者和消費者的原理。

2.1 生產(chǎn)端問題

生產(chǎn)者原理

Kafka生產(chǎn)者生產(chǎn)消息后,會將消息發(fā)送到Kafka集群的Leader中,然后Kafka集群的Leader收到消息后會返回ACK確認消息給生產(chǎn)者Producer。主要拆解為以下幾個步驟。

  • Producer先從Kafka集群找到該Partition的Leader。
  • Producer將消息發(fā)送給Leader,Leader將該消息寫入本地。
  • Follwer從Leader pull消息,寫入本地Log后Leader發(fā)送ACK。
  • Leader 收到所有 ISR 中的 Replica 的 ACK 后,增加High Watermark,并向 Producer 發(fā)送 ACK。

  • 因此,Kafka集群(其實是分區(qū)的Leader)最終會返回一個ACK來確認Producer推送消息的結(jié)果,這里Kafka提供了三種模式:
  • NoResponse RequiredAcks = 0:這個代表的就是不進行消息推送是否成功的確認。
  • WaitForLocal RequiredAcks = 1:當local(Leader)確認接收成功后,就可以返回了。
  • WaitForAll RequiredAcks = -1:當所有的Leader和Follower都接收成功時,才會返回。

因此這個配置的影響也分為下面三種情況:

  • 設置為0,Producer不進行消息發(fā)送的確認,Kafka集群(Broker)可能由于一些原因并沒有收到對應消息,從而引起消息丟失。
  • 設置為1,Producer在確認到 Topic Leader 已經(jīng)接收到消息后,完成發(fā)送,此時有可能 Follower 并沒有接收到對應消息。此時如果 Leader 突然宕機,在經(jīng)過選舉之后,沒有接到消息的 Follower 晉升為 Leader,從而引起消息丟失。
  • 設置為-1,可以很好的確認Kafka集群是否已經(jīng)完成消息的接收和本地化存儲,并且可以在Producer發(fā)送失敗時進行重試。

生產(chǎn)端解決消息丟失方案:

  • 通過設置RequiredAcks模式來解決,選用WaitForAll(對應值為-1)可以保證數(shù)據(jù)推送成功,不過會影響延時。
  • 引入重試機制,設置重試次數(shù)和重試間隔。
  • 當然,最后就是使用Kafka的多副本機制保證Kafka集群本身的可靠性,確保當Leader掛掉之后能進行Follower選舉晉升為新的Leader。

2.2 消費端問題

消費端的消息丟失問題
消費端的消息丟失主要是因為在消費過程中出現(xiàn)了異常,但是對應消息的 Offset 已經(jīng)提交,那么消費異常的消息將會丟失。
前面介紹過,Offset的提交包括手動提交和自動提交,可通過kafka.consumer.enable-auto-commit進行配置。
手動提交可以靈活的確認是否將本次消費數(shù)據(jù)的Offset進行提交,可以很好的避免消息丟失的情況。
自動提交是引起消息丟失的主要誘因。因為消息的消費并不會影響到Offset的提交。
大部分的解決方案為了盡可能的保證數(shù)據(jù)的完整性,都是盡量去選用手動提交的方式,當數(shù)據(jù)處理完之后再進行提交。
當然,在golang中我們主要使用sarama包的Kafka,sarama自動提交的原理是先進行標記,再進行提交,如下代碼所示:

type exampleConsumerGroupHandler struct{}
func (exampleConsumerGroupHandler) Setup(_ ConsumerGroupSession) error   { return nil }
func (exampleConsumerGroupHandler) Cleanup(_ ConsumerGroupSession) error { return nil }
func (h exampleConsumerGroupHandler) ConsumeClaim(sess ConsumerGroupSession, claim ConsumerGroupClaim) error {
   for msg := range claim.Messages() {
      fmt.Printf("Message topic:%q partition:%d offset:%d
", msg.Topic, msg.Partition, msg.Offset)
      // 標記消息已處理,sarama會自動提交
      // 處理數(shù)據(jù)(如真正持久化mysql...)
      sess.MarkMessage(msg, "")
   }
   return nil

因此,我們完全可以在標記之前進行數(shù)據(jù)的處理,例如插入Mysql等,當出現(xiàn)插入成功后程序崩潰,下一次最多重復消費一次(因為還沒標記,Offset沒有提交),而不會因為Offset超前,導致應用層消息丟失了。

手動提交模式下當然是很靈活的控制的,但確實已經(jīng)沒必要了:

consumerConfig := sarama.NewConfig()
consumerConfig.Version = sarama.V2_8_0_0
consumerConfig.Consumer.Return.Errors = false
consumerConfig.Consumer.Offsets.AutoCommit.Enable = false  // 禁用自動提交,改為手動
consumerConfig.Consumer.Offsets.Initial = sarama.OffsetNewest
func (h msgConsumerGroup) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
   for msg := range claim.Messages() {
      fmt.Printf("%s Message topic:%q partition:%d offset:%d  value:%s
", h.name, msg.Topic, msg.Partition, msg.Offset, string(msg.Value))
      // 插入mysql....
      // 手動提交模式下,也需要先進行標記
      sess.MarkMessage(msg, "")
      consumerCount++
      if consumerCount%3 == 0 {
         // 手動提交,不能頻繁調(diào)用
         t1 := time.Now().Nanosecond()
         sess.Commit()
         t2 := time.Now().Nanosecond()
         fmt.Println("commit cost:", (t2-t1)/(1000*1000), "ms")
      }
   }
   return nil
}

到此這篇關(guān)于Golang中Kafka的重復消費和消息丟失問題的解決方案的文章就介紹到這了,更多相關(guān)Golang Kafka重復消費和消息丟失內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • go語言編程之select信道處理示例詳解

    go語言編程之select信道處理示例詳解

    這篇文章主要為大家介紹了go語言編程之select信道處理示例詳解,<BR>有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步早日升職加薪
    2022-04-04
  • 一文詳解kubernetes?中資源分配的那些事

    一文詳解kubernetes?中資源分配的那些事

    這篇文章主要為大家介紹了kubernetes?中資源分配的那些事,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-04-04
  • go語言執(zhí)行等待直到后臺goroutine執(zhí)行完成實例分析

    go語言執(zhí)行等待直到后臺goroutine執(zhí)行完成實例分析

    這篇文章主要介紹了go語言執(zhí)行等待直到后臺goroutine執(zhí)行完成的方法,實例分析了Go語言中WaitGroup的使用技巧,需要的朋友可以參考下
    2015-03-03
  • Go語言defer語句的三種機制整理

    Go語言defer語句的三種機制整理

    在本篇文章里小編給大家分享的是一篇關(guān)于Go語言defer語句的三種機制整理,需要的朋友們學習下吧。
    2020-03-03
  • 基于Go語言實現(xiàn)一個并發(fā)端口掃描器

    基于Go語言實現(xiàn)一個并發(fā)端口掃描器

    這篇文章主要介紹了如何使用 Go 實現(xiàn)一個并發(fā)端口掃描器,通過 Goroutine 并發(fā)掃描多個端口,極大地提升端口掃描的效率,本文不僅講解了如何使用 Go 的并發(fā)特性,還涉及了如何處理超時和錯誤,保證端口掃描的健壯性和效率,需要的朋友可以參考下
    2025-08-08
  • Go語言學習之文件操作方法詳解

    Go語言學習之文件操作方法詳解

    這篇文章主要為大家詳細介紹了Go語言中一些常見的文件操作,文中的示例代碼講解詳細,對我們學習Go語言有一定的幫助,需要的可以參考一下
    2022-04-04
  • golang內(nèi)存逃逸分析

    golang內(nèi)存逃逸分析

    本文主要介紹了golang內(nèi)存逃逸分析,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2025-06-06
  • Golang之sync.Pool對象池對象重用機制總結(jié)

    Golang之sync.Pool對象池對象重用機制總結(jié)

    這篇文章主要對Golang的sync.Pool對象池對象重用機制做了一個總結(jié),文中有相關(guān)的代碼示例和圖解,具有一定的參考價值,需要的朋友可以參考下
    2023-07-07
  • golang數(shù)組和切片作為參數(shù)和返回值的實現(xiàn)

    golang數(shù)組和切片作為參數(shù)和返回值的實現(xiàn)

    本文主要介紹了golang數(shù)組和切片作為參數(shù)和返回值的實現(xiàn),文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-02-02
  • Golang利用compress/flate包來壓縮和解壓數(shù)據(jù)

    Golang利用compress/flate包來壓縮和解壓數(shù)據(jù)

    在處理需要高效存儲和快速傳輸?shù)臄?shù)據(jù)時,數(shù)據(jù)壓縮成為了一項不可或缺的技術(shù),Go語言的compress/flate包為我們提供了對DEFLATE壓縮格式的原生支持,本文將深入探討compress/flate包的使用方法,揭示如何利用它來壓縮和解壓數(shù)據(jù),并提供實際的代碼示例,需要的朋友可以參考下
    2024-08-08

最新評論

余干县| 张家界市| 鄂尔多斯市| 留坝县| 呼玛县| 盱眙县| 黑河市| 分宜县| 寻乌县| 九台市| 和田市| 太白县| 昭通市| 彭山县| 泗水县| 奉节县| 高陵县| 金寨县| 论坛| 平远县| 莱芜市| 桐乡市| 托里县| 南丰县| 沧源| 罗田县| 唐山市| 郓城县| 皋兰县| 平果县| 买车| 临澧县| 定兴县| 佛山市| 甘德县| 余庆县| 桐乡市| 南乐县| 武隆县| 留坝县| 东港市|