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

Go+Kafka實現(xiàn)延遲消息的實現(xiàn)示例

 更新時間:2022年07月25日 08:29:15   作者:jiaxwu  
本文主要介紹了Go+Kafka實現(xiàn)延遲消息的實現(xiàn)示例,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧

前言

延遲隊列是一個非常有用的工具,我們經(jīng)常遇到需要使用延遲隊列的場景,比如延遲通知,訂單關(guān)閉等等。

這篇文章主要是使用Go+Kafka實現(xiàn)延遲消息。

使用了sarama客戶端。

原理

Kafka實現(xiàn)延遲消息分為下面三步:

  • 生產(chǎn)者把消息發(fā)送到延遲隊列
  • 延遲服務(wù)把延遲隊列里超過延遲時間的消息寫入真實隊列
  • 消費者消費真實隊列里的消息

簡單的實現(xiàn)

生產(chǎn)者

生產(chǎn)者只是把消息發(fā)送到延遲隊列

msg := &sarama.ProducerMessage{
   Topic: kafka_delay_queue_test.DelayTopic,
   Value: sarama.ByteEncoder("test" + strconv.Itoa(i)),
}
if _, _, err := producer.SendMessage(msg); err != nil {
   log.Println(err)
}

延遲服務(wù)

延遲服務(wù)會訂閱延遲隊列的消息,并把超時消息發(fā)送到真實隊列

if err = consumerGroup.Consume(context.Background(),
   []string{kafka_delay_queue_test.DelayTopic}, consumer); err != nil {
   break
}
type Consumer struct {
   producer sarama.SyncProducer
   delay    time.Duration
}

func NewConsumer(producer sarama.SyncProducer, delay time.Duration) *Consumer {
   return &Consumer{
      producer: producer,
      delay:    delay,
   }
}

func (c *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
   for message := range claim.Messages() {
      // 如果消息已經(jīng)超時,把消息發(fā)送到真實隊列
      now := time.Now()
      if now.Sub(message.Timestamp) >= c.delay {
         _, _, err := c.producer.SendMessage(&sarama.ProducerMessage{
            Topic: kafka_delay_queue_test.RealTopic,
            Key:   sarama.ByteEncoder(message.Key),
            Value: sarama.ByteEncoder(message.Value),
         })
         if err == nil {
            session.MarkMessage(message, "")
         }
         continue
      }
      // 否則休眠一秒
      time.Sleep(time.Second)
      return nil
   }
   return nil
}

消費者

消費者只是訂閱真實隊列并消費消息

if err = consumerGroup.Consume(context.Background(), 
   []string{kafka_delay_queue_test.RealTopic}, consumer); err != nil {
   break
}
type Consumer struct{}

func NewConsumer() *Consumer {
   return &Consumer{}
}

func (c *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
   for message := range claim.Messages() {
      fmt.Println("收到消息:", message.Value, message.Timestamp)
      session.MarkMessage(message, "")
   }
   return nil
}

改進(jìn)點

通用的延遲服務(wù)

可以把延遲服務(wù)封裝成一個通用的服務(wù),這樣生產(chǎn)者可以直接把消息發(fā)送給延遲服務(wù),讓延遲服務(wù)去處理剩下的邏輯。

延遲服務(wù)可以提供多個延時等級,比如5s、10s、30s、1m、5m、10m、1h、2h等,類似于RocketMQ。

生產(chǎn)者負(fù)責(zé)延遲服務(wù)

也可以讓生產(chǎn)者負(fù)責(zé)延遲服務(wù),讓生產(chǎn)者自己把延遲隊列里面的消息發(fā)送到真實隊列。

下面是一個簡單的實現(xiàn):

// KafkaDelayQueueProducer 延遲隊列生產(chǎn)者,包含了生產(chǎn)者和延遲服務(wù)
type KafkaDelayQueueProducer struct {
   producer   sarama.SyncProducer // 生產(chǎn)者
   delayTopic string              // 延遲服務(wù)主題
}

// NewKafkaDelayQueueProducer 創(chuàng)建延遲隊列生產(chǎn)者
// producer 生產(chǎn)者
// delayServiceConsumerGroup 延遲服務(wù)消費者
// delayTime 延遲時間
// delayTopic 延遲服務(wù)主題
// realTopic 真實隊列主題
func NewKafkaDelayQueueProducer(producer sarama.SyncProducer, delayServiceConsumerGroup sarama.ConsumerGroup,
   delayTime time.Duration, delayTopic, realTopic string) *KafkaDelayQueueProducer {
   // 啟動延遲服務(wù)
   consumer := NewDelayServiceConsumer(producer, delayTime, realTopic)
   go func() {
      for {
         if err := delayServiceConsumerGroup.Consume(context.Background(),
            []string{delayTopic}, consumer); err != nil {
            break
         }
      }
   }()
   return &KafkaDelayQueueProducer{
      producer:   producer,
      delayTopic: delayTopic,
   }
}

// SendMessage 發(fā)送消息
func (q *KafkaDelayQueueProducer) SendMessage(msg *sarama.ProducerMessage) (partition int32, offset int64, err error) {
   msg.Topic = q.delayTopic
   return q.producer.SendMessage(msg)
}

// DelayServiceConsumer 延遲服務(wù)消費者
type DelayServiceConsumer struct {
   producer  sarama.SyncProducer
   delay     time.Duration
   realTopic string
}

func NewDelayServiceConsumer(producer sarama.SyncProducer, delay time.Duration,
   realTopic string) *DelayServiceConsumer {
   return &DelayServiceConsumer{
      producer:  producer,
      delay:     delay,
      realTopic: realTopic,
   }
}

func (c *DelayServiceConsumer) ConsumeClaim(session sarama.ConsumerGroupSession,
   claim sarama.ConsumerGroupClaim) error {
   for message := range claim.Messages() {
      // 如果消息已經(jīng)超時,把消息發(fā)送到真實隊列
      now := time.Now()
      if now.Sub(message.Timestamp) >= c.delay {
         _, _, err := c.producer.SendMessage(&sarama.ProducerMessage{
            Topic: c.realTopic,
            Key:   sarama.ByteEncoder(message.Key),
            Value: sarama.ByteEncoder(message.Value),
         })
         if err == nil {
            session.MarkMessage(message, "")
         }
         continue
      }
      // 否則休眠一秒
      time.Sleep(time.Second)
      return nil
   }
   return nil
}

func (c *DelayServiceConsumer) Setup(sarama.ConsumerGroupSession) error {
   return nil
}

func (c *DelayServiceConsumer) Cleanup(sarama.ConsumerGroupSession) error {
   return nil
}

總結(jié)

使用中間隊列+輪詢可以很容易的在Kafka實現(xiàn)延遲消息,如果需要一個通用的延遲隊列也可以實現(xiàn)一個通用的延遲服務(wù),也可以讓消費者負(fù)責(zé)延遲服務(wù)的功能。

完整代碼:

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

相關(guān)文章

  • Go語言中一些不常見的命令參數(shù)詳解

    Go語言中一些不常見的命令參數(shù)詳解

    這篇文章主要給大家介紹了關(guān)于Go語言中一些不常見的命令參數(shù)的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧。
    2017-12-12
  • Go語言實現(xiàn)彩色輸出示例詳解

    Go語言實現(xiàn)彩色輸出示例詳解

    這篇文章主要為大家介紹了Go語言實現(xiàn)彩色輸出示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-09-09
  • Golang并發(fā)發(fā)送HTTP請求的各種方法

    Golang并發(fā)發(fā)送HTTP請求的各種方法

    在 Golang 領(lǐng)域,并發(fā)發(fā)送 HTTP 請求是優(yōu)化 Web 應(yīng)用程序的一項重要技能,本文探討了實現(xiàn)此目的的各種方法,從基本的 goroutine 到涉及通道和sync.WaitGroup 的高級技術(shù),需要的朋友可以參考下
    2024-02-02
  • go使用errors.Wrapf()代替log.Error()方法示例

    go使用errors.Wrapf()代替log.Error()方法示例

    這篇文章主要為大家介紹了go使用errors.Wrapf()代替log.Error()的方法示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-08-08
  • Goland使用delve進(jìn)行遠(yuǎn)程調(diào)試的詳細(xì)教程

    Goland使用delve進(jìn)行遠(yuǎn)程調(diào)試的詳細(xì)教程

    網(wǎng)上給出的使用delve進(jìn)行遠(yuǎn)程調(diào)試,都需要先在本地交叉編譯或者在遠(yuǎn)程主機上編譯出可運行的程序,然后再用delve在遠(yuǎn)程啟動程序,本教程會將上面的步驟簡化為只需要兩步,1,在遠(yuǎn)程運行程序2,在本地啟動調(diào)試,需要的朋友可以參考下
    2024-08-08
  • Go?中?time.After?可能導(dǎo)致的內(nèi)存泄露問題解析

    Go?中?time.After?可能導(dǎo)致的內(nèi)存泄露問題解析

    這篇文章主要介紹了Go?中?time.After?可能導(dǎo)致的內(nèi)存泄露,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2023-05-05
  • Golang限流庫與漏桶和令牌桶的使用介紹

    Golang限流庫與漏桶和令牌桶的使用介紹

    這篇文章主要介紹了golang限流庫以及漏桶與令牌桶的實現(xiàn)原理,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧
    2023-03-03
  • 解決golang時間字符串轉(zhuǎn)time.Time的坑

    解決golang時間字符串轉(zhuǎn)time.Time的坑

    這篇文章主要介紹了解決golang時間字符串轉(zhuǎn)time.Time的坑,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-04-04
  • Go語言數(shù)據(jù)類型詳細(xì)介紹

    Go語言數(shù)據(jù)類型詳細(xì)介紹

    這篇文章主要介紹了Go語言數(shù)據(jù)類型詳細(xì)介紹,Go語言數(shù)據(jù)類型包含基礎(chǔ)類型和復(fù)合類型兩大類,下文關(guān)于這兩類型的相關(guān)介紹,需要的小伙伴可以參考一下
    2022-03-03
  • 1行Go代碼實現(xiàn)反向代理的示例

    1行Go代碼實現(xiàn)反向代理的示例

    這篇文章主要介紹了1行Go代碼實現(xiàn)反向代理的示例,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-08-08

最新評論

马龙县| 长丰县| 三河市| 腾冲县| 桐乡市| 凤山县| 达州市| 大方县| 霍林郭勒市| 陆川县| 许昌市| 原平市| 南澳县| 新昌县| 敦化市| 体育| 四川省| 福安市| 宜黄县| 象州县| 鲁山县| 永泰县| 松阳县| 延寿县| 海城市| 申扎县| 九寨沟县| 客服| 吴旗县| 怀宁县| 德阳市| 友谊县| 平谷区| 砚山县| 衡山县| 繁峙县| 固原市| 长汀县| 鲁山县| 砀山县| 老河口市|