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

Go+Redis實現(xiàn)延遲隊列實操

 更新時間:2022年09月14日 08:54:28   作者:jiaxwu???????  
這篇文章主要介紹了Go+Redis實現(xiàn)延遲隊列實操,延遲隊列是一種非常使用的數(shù)據(jù)結(jié)構(gòu),我們經(jīng)常有需要延遲推送處理消息的場景,比如延遲60秒發(fā)送短信,延遲30分鐘關(guān)閉訂單,消息消費失敗延遲重試等

前言

延遲隊列是一種非常使用的數(shù)據(jù)結(jié)構(gòu),我們經(jīng)常有需要延遲推送處理消息的場景,比如延遲60秒發(fā)送短信,延遲30分鐘關(guān)閉訂單,消息消費失敗延遲重試等等。

一般我們實現(xiàn)延遲消息都需要依賴底層的有序結(jié)構(gòu),比如堆,而Redis剛好提供了zset這種數(shù)據(jù)類型,它的底層實現(xiàn)是哈希表+跳表,也是一種有序的結(jié)構(gòu),所以這篇文章主要是使用Go+Redis來實現(xiàn)延遲隊列。

當(dāng)然Redis本身并不支持延遲隊列,所以我們只是實現(xiàn)一個比較簡單的延遲隊列,而且Redis不太適合大量消息堆積,所以只適合比較簡單的場景,如果需要更加強大穩(wěn)定的消息隊列,可以使用RocketMQ等自帶延遲消息的消息隊列。

我們這里先定一下我們要實現(xiàn)的幾個目標(biāo):

  • 消息必須至少被消費一次
  • 多個生產(chǎn)者
  • 多個消費者

然后我們定義一個簡單的接口:

  • Push(msg) error:添加消息到隊列
  • Consume(topic, batchSize, func(msg) error):消費消息

簡單的實現(xiàn)

  • 每個主題最多可以被一個消費者消費,因為不會對主題進(jìn)行分區(qū)
  • 但是可以多個生產(chǎn)者同時進(jìn)行生產(chǎn),因為Push操作是原子的
  • 同時需要消費操作返回值error為nil才刪除消息,保證消息至少被消費一次

定義消息

這個消息參考了Kafka的消息結(jié)構(gòu):

  • Topic可以是某個隊列的名字
  • Key是消息的唯一標(biāo)識,在一個隊列里面不可以重復(fù)
  • Body是消息的內(nèi)容
  • Delay是消息的延遲時間
  • ReadyTime是消息準(zhǔn)備好執(zhí)行的時間
// Msg 消息
type Msg struct {
   Topic     string        // 消息的主題
   Key       string        // 消息的Key
   Body      []byte        // 消息的Body
   Delay     time.Duration // 延遲時間(秒)
   ReadyTime time.Time     // 消息準(zhǔn)備好執(zhí)行的時間(now + delay)
}

Push

由于我們需要把消息的Body存儲到Hash,把消息的ReadyTime存儲到ZSet,所以我們需要一個簡單的Lua腳本來保證這兩個操作是原子的。

同時我們不會覆蓋已經(jīng)存在的相同Key的消息。

const delayQueuePushRedisScript = `
-- KEYS[1]: topicZSet
-- KEYS[2]: topicHash
-- ARGV[1]: 消息的Key
-- ARGV[2]: 消息的Body
-- ARGV[3]: 消息準(zhǔn)備好執(zhí)行的時間

local topicZSet = KEYS[1]
local topicHash = KEYS[2]
local key = ARGV[1]
local body = ARGV[2]
local readyTime = tonumber(ARGV[3])

-- 添加readyTime到zset
local count = redis.call("zadd", topicZSet, readyTime, key)
-- 消息已經(jīng)存在
if count == 0 then 
   return 0
end
-- 添加body到hash
redis.call("hsetnx", topicHash, key, body)
return 1
`
func (q *SimpleRedisDelayQueue) Push(ctx context.Context, msg *Msg) error {
   // 如果設(shè)置了ReadyTime,就使用RedisTime
   var readyTime int64
   if !msg.ReadyTime.IsZero() {
      readyTime = msg.ReadyTime.Unix()
   } else {
      // 否則使用Delay
      readyTime = time.Now().Add(msg.Delay).Unix()
   }
   success, err := q.pushScript.Run(ctx, q.client, []string{q.topicZSet(msg.Topic), q.topicHash(msg.Topic)},
      msg.Key, msg.Body, readyTime).Bool()
   if err != nil {
      return err
   }
   if !success {
      return ErrDuplicateMessage
   }
   return nil
}

Consume

其中第二個參數(shù)batchSize表示用于批量獲取已經(jīng)準(zhǔn)備好執(zhí)行的消息,減少網(wǎng)絡(luò)請求。

fn是對消息進(jìn)行處理的函數(shù),它有一個返回值error,如果是nil才表示消息消費成功,然后調(diào)用刪除腳本把成功消費的消息給刪除(需要原子的刪除ZSet和Hash里面的內(nèi)容)。

const delayQueueDelRedisScript = `
-- KEYS[1]: topicZSet
-- KEYS[2]: topicHash
-- ARGV[1]: 消息的Key

local topicZSet = KEYS[1]
local topicHash = KEYS[2]
local key = ARGV[1]

-- 刪除zset和hash關(guān)于這條消息的內(nèi)容
redis.call("zrem", topicZSet, key)
redis.call("hdel", topicHash, key)
return 1
`
func (q *SimpleRedisDelayQueue) Consume(topic string, batchSize int, fn func(msg *Msg) error) {
   for {
      // 批量獲取已經(jīng)準(zhǔn)備好執(zhí)行的消息
      now := time.Now().Unix()
      zs, err := q.client.ZRangeByScoreWithScores(context.Background(), q.topicZSet(topic), &redis.ZRangeBy{
         Min:   "-inf",
         Max:   strconv.Itoa(int(now)),
         Count: int64(batchSize),
      }).Result()
      // 如果獲取出錯或者獲取不到消息,則休眠一秒
      if err != nil || len(zs) == 0 {
         time.Sleep(time.Second)
         continue
      }
      // 遍歷每個消息
      for _, z := range zs {
         key := z.Member.(string)
         // 獲取消息的body
         body, err := q.client.HGet(context.Background(), q.topicHash(topic), key).Bytes()
         if err != nil {
            continue
         }

         // 處理消息
         err = fn(&Msg{
            Topic:     topic,
            Key:       key,
            Body:      body,
            ReadyTime: time.Unix(int64(z.Score), 0),
         })
         if err != nil {
            continue
         }

         // 如果消息處理成功,刪除消息
         q.delScript.Run(context.Background(), q.client, []string{q.topicZSet(topic), q.topicHash(topic)}, key)
      }
   }
}

存在的問題

如果多個線程同時調(diào)用Consume函數(shù),那么多個線程會拉取相同的可執(zhí)行的消息,造成消息重復(fù)的被消費。

多消費者實現(xiàn)

  • 每個主題最多可以被分區(qū)個數(shù)個消費者消費,會對主題進(jìn)行分區(qū)

定義消息

  • 我們添加了一個Partition字段表示消息的分區(qū)號
// Msg 消息
type Msg struct {
   Topic     string        // 消息的主題
   Key       string        // 消息的Key
   Body      []byte        // 消息的Body
   Partition int           // 分區(qū)號
   Delay     time.Duration // 延遲時間(秒)
   ReadyTime time.Time     // 消息準(zhǔn)備好執(zhí)行的時間
}

Push

代碼與SimpleRedisDelayQueue的Push相同,只是我們會使用Msg里面的Partition字段對主題進(jìn)行分區(qū)。

func (q *PartitionRedisDelayQueue) topicZSet(topic string, partition int) string {
   return fmt.Sprintf("%s:%d:z", topic, partition)
}

func (q *PartitionRedisDelayQueue) topicHash(topic string, partition int) string {
   return fmt.Sprintf("%s:%d:h", topic, partition)
}

Consume

代碼與SimpleRedisDelayQueue的Consume相同,我們也只是對Consume多加了一個partition參數(shù)用于指定消費的分區(qū)。

func (q *PartitionRedisDelayQueue) Consume(topic string, batchSize, partition int, fn func(msg *Msg) error) {
    // ...
}

存在的問題

一個比較大的問題就是我們需要手動指定分區(qū)而不是自動分配分區(qū),這個問題對于Push操作解決起來比較容易,可以通過哈希算法對Key進(jìn)行哈希取模進(jìn)行分區(qū),比如murmur3。但是對于Consume就比較復(fù)雜,因為我們必須記錄哪個分區(qū)已經(jīng)被消費者消費了。如果真的需要更加復(fù)雜的場景還是建議使用RocketMQ、Kafka等消息隊列進(jìn)行實現(xiàn)。

總結(jié)

  • 使用Redis的ZSet可以很容易的實現(xiàn)一個高性能消息隊列
  • 但是Redis的ZSet實現(xiàn)的消息隊列不適合大量消息堆積的場景,同時如果需要實現(xiàn)自動分區(qū)消費功能會比較復(fù)雜
  • 適合消息量不是很大,且不是很復(fù)雜的場景
  • 如果需要大量堆積消息和穩(wěn)定的多消費者功能,可以使用自帶延遲消息的RocketMQ

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

相關(guān)文章

  • golang正則之命名分組方式

    golang正則之命名分組方式

    這篇文章主要介紹了golang正則之命名分組方式,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-04-04
  • 詳解如何在Go中如何編寫出可測試的代碼

    詳解如何在Go中如何編寫出可測試的代碼

    在編寫測試代碼之前,還有一個很重要的點,容易被忽略,就是什么樣的代碼是可測試的代碼,所以本文就來聊一聊在?Go?中如何寫出可測試的代碼吧
    2023-08-08
  • Golang編程并發(fā)工具庫MapReduce使用實踐

    Golang編程并發(fā)工具庫MapReduce使用實踐

    這篇文章主要為大家介紹了Golang并發(fā)工具庫MapReduce的使用實踐,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-04-04
  • Go語言并發(fā)編程基礎(chǔ)上下文概念詳解

    Go語言并發(fā)編程基礎(chǔ)上下文概念詳解

    這篇文章主要為大家介紹了Go語言并發(fā)編程基礎(chǔ)上下文示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-08-08
  • 如何有效控制Go線程數(shù)實例探究

    如何有效控制Go線程數(shù)實例探究

    這篇文章主要為大家介紹了如何有效控制?Go?線程數(shù)的問題探究,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2024-01-01
  • Go語言使用Timeout Context取消任務(wù)的實現(xiàn)

    Go語言使用Timeout Context取消任務(wù)的實現(xiàn)

    本文主要介紹了Go語言使用Timeout Context取消任務(wù)的實現(xiàn),包括基本的任務(wù)取消和控制HTTP客戶端請求的超時,具有一定的參考價值,感興趣的可以了解一下
    2024-01-01
  • go xorm存庫處理null值問題

    go xorm存庫處理null值問題

    這篇文章主要介紹了go xorm存庫處理null值問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2023-12-12
  • Golang的循環(huán)語句和循環(huán)控制語句詳解

    Golang的循環(huán)語句和循環(huán)控制語句詳解

    循環(huán)語句為了簡化程序中有規(guī)律的重復(fù)性操作,需要用到循環(huán)語句,和其他大多數(shù)編程語言一樣,GO的循環(huán)語句有for循環(huán),不同的是沒有while循環(huán),而循環(huán)控制語句可以改變循環(huán)語句的執(zhí)行過程,下面給大家介紹下go循環(huán)語句和循環(huán)控制語句的相關(guān)知識,一起看看吧
    2021-11-11
  • Go語言實現(xiàn)猜數(shù)字小游戲

    Go語言實現(xiàn)猜數(shù)字小游戲

    這篇文章主要為大家詳細(xì)介紹了Go語言實現(xiàn)猜數(shù)字小游戲,文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2020-10-10
  • golang常見接口限流算法的實現(xiàn)

    golang常見接口限流算法的實現(xiàn)

    本文主要介紹了golang常見接口限流算法的實現(xiàn),包含固定窗口、滑動窗口、漏桶和令牌桶,具有一定的參考價值,感興趣的可以了解一下
    2025-03-03

最新評論

修武县| 保靖县| 永登县| 高唐县| 奉节县| 舞钢市| 多伦县| 双柏县| 邵阳市| 邵阳市| 灵寿县| 竹山县| 满洲里市| 博乐市| 长岛县| 鲁甸县| 家居| 安阳市| 沂源县| 观塘区| 建平县| 阜南县| 绥棱县| 光山县| 泾川县| 四川省| 获嘉县| 洪雅县| 永和县| 天镇县| 布尔津县| 会泽县| 易门县| 壤塘县| 诏安县| 西城区| 敦煌市| 当阳市| 河南省| 库伦旗| 宜良县|