中文字幕av专区_日韩电影在线播放_精品国产精品久久一区免费式_av在线免费观看网站

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

golang延時任務如何實現

發布時間:2023-03-21 14:42:55 來源:億速云 閱讀:115 作者:iii 欄目:開發技術

這篇文章主要介紹“golang延時任務如何實現”,在日常操作中,相信很多人在golang延時任務如何實現問題上存在疑惑,小編查閱了各式資料,整理出簡單好用的操作方法,希望對大家解答”golang延時任務如何實現”的疑惑有所幫助!接下來,請跟著小編一起來學習吧!

    實現思路

    我們都知道,任何一種隊列,實際上都是存在生產者和消費者兩部分的。只不過,延時任務相對于普通隊列,多了一個延時的特性罷了。

    1、生產者

    從生產者的角度上講,當用戶推送一個任務過來的時候,會攜帶著延遲執行的時間數值。為了讓這個任務到預定時刻能執行,我們需要將這個任務放在內存里儲存一段時間,并且時間是一維的,在不斷增長。那么,我們用什么數據結構存儲呢?

    (1)選擇一:map。由于map具有無序性,無法按照執行時間排序,我們無法保證取出的任務是否是當前時間點需要執行的,所以排除這個選項。

    (2)選擇二:channel。的確,channel有時候可以看作隊列,然而,它的輸出和輸入嚴格遵循著“先進先出”的原則,遺憾的是,先進的任務未必就是先執行的,因此,channel也并不合適。

    (3)選擇三:slice。切片貌似可行,因為切片元素是具有有序性的,所以,如果我們能夠按照執行時間的順序排列好所有的切片元素,那么,每次只要讀取切片的頭元素(也可能是尾元素),就可以得到我們要的任務。

    2、消費者

    從消費者的角度來說,它最大的難點在于,如何讓每個任務,在特定的時間點被消費。那么,針對每一個任務,我們如何實現,讓它等待一段時間后再執行呢?

    沒錯,就是timer。

    總結下來,“切片+timer”的組合,應該是可以達到目的的。

    步步為營

    1、數據流

    (1)用戶調用InitDelayQueue() ,初始化延時任務對象。

    (2)開啟協程,監聽任務操作管道(add/delete信號),以及執行時間管道(timer.C信號)。

    (3)用戶發出add/delete信號。

    (4)(2)中的協程捕捉到(3)中的信號,對任務列表進行變更。

    (5)當任務執行的時間點到達的時候(timer.C管道有元素輸出的時候),執行任務。

    golang延時任務如何實現

    2、數據結構

    (1)延時任務對象

    // 延時任務對象
    type DelayQueue struct {
       tasks                 []*task             // 存儲任務列表的切片
       add                   chan *task          // 用戶添加任務的管道信號
       remove                chan string         // 用戶刪除任務的管道信號
       waitRemoveTaskMapping map[string]struct{} // 等待刪除的任務id列表
    }

    這里需要注意,有一個waitRemoveTaskMapping字段。由于要刪除的任務,可能還在add管道中,沒有及時更新到tasks字段中,所以,需要臨時記錄下客戶要刪除的任務id。

    (2)任務對象

    // 任務對象
    type task struct {
       id       string    // 任務id
       execTime time.Time // 執行時間
       f        func()    // 執行函數
    }
    3、初始化延時任務對象
    // 初始化延時任務對象
    func InitDelayQueue() *DelayQueue {
       q := &DelayQueue{
          add:                   make(chan *task, 10000),
          remove:                make(chan string, 100),
          waitRemoveTaskMapping: make(map[string]struct{}),
       }
       return q
    }

    在這個過程中,我們需要對用戶對任務的操作信號,以及任務的執行時間信號進行監聽。

    func (q *DelayQueue) start() {
       for {
          // to do something...
          select {
          case now := <-timer.C:
             // 任務執行時間信號
             // to do something...
          case t := <-q.add:
             // 任務推送信號
             // to do something...
          case id := <-q.remove:
             // 任務刪除信號
             // to do something...
          }
       }
    }

    完善我們的初始化方法:

    // 初始化延時任務對象
    func InitDelayQueue() *DelayQueue {
       q := &DelayQueue{
          add:                   make(chan *task, 10000),
          remove:                make(chan string, 100),
          waitRemoveTaskMapping: make(map[string]struct{}),
       }
       // 開啟協程,監聽任務相關信號
       go q.start()
       return q
    }
    4、生產者推送任務

    生產者推送任務的時候,只需要將任務加到add管道中即可,在這里,我們生成一個任務id,并返回給用戶。

    // 用戶推送任務
    func (q *DelayQueue) Push(timeInterval time.Duration, f func()) string {
       // 生成一個任務id,方便刪除使用
       id := genTaskId()
       t := &task{
          id:       id,
          execTime: time.Now().Add(timeInterval),
          f:        f,
       }
       // 將任務推到add管道中
       q.add <- t
       return id
    }
    5、任務推送信號的處理

    在這里,我們要將用戶推送的任務放到延時任務的tasks字段中。由于,我們需要將任務按照執行時間順序排序,所以,我們需要找到新增任務在切片中的插入位置。又因為,插入之前的任務列表已經是有序的,所以,我們可以采用二分法處理。

    // 使用二分法判斷新增任務的插入位置
    func (q *DelayQueue) getTaskInsertIndex(t *task, leftIndex, rightIndex int) (index int) {
       if len(q.tasks) == 0 {
          return
       }
       length := rightIndex - leftIndex
       if q.tasks[leftIndex].execTime.Sub(t.execTime) >= 0 {
          // 如果當前切片中最小的元素都超過了插入的優先級,則插入位置應該是最左邊
          return leftIndex
       }
       if q.tasks[rightIndex].execTime.Sub(t.execTime) <= 0 {
          // 如果當前切片中最大的元素都沒超過插入的優先級,則插入位置應該是最右邊
          return rightIndex + 1
       }
       if length == 1 && q.tasks[leftIndex].execTime.Before(t.execTime) && q.tasks[rightIndex].execTime.Sub(t.execTime) >= 0 {
          // 如果插入的優先級剛好在僅有的兩個優先級之間,則中間的位置就是插入位置
          return leftIndex + 1
       }
       middleVal := q.tasks[leftIndex+length/2].execTime
       // 這里用二分法遞歸的方式,一直尋找正確的插入位置
       if t.execTime.Sub(middleVal) <= 0 {
          return q.getTaskInsertIndex(t, leftIndex, leftIndex+length/2)
       } else {
          return q.getTaskInsertIndex(t, leftIndex+length/2, rightIndex)
       }
    }

    找到正確的插入位置后,我們才能將任務準確插入:

    // 將任務添加到任務切片列表中
    func (q *DelayQueue) addTask(t *task) {
       // 尋找新增任務的插入位置
       insertIndex := q.getTaskInsertIndex(t, 0, len(q.tasks)-1)
       // 找到了插入位置,更新任務列表
       q.tasks = append(q.tasks, &task{})
       copy(q.tasks[insertIndex+1:], q.tasks[insertIndex:])
       q.tasks[insertIndex] = t
    }

    那么,在監聽add管道的時候,我們直接調用上述addTask() 即可。

    func (q *DelayQueue) start() {
       for {
          // to do something...
          select {
          case now := <-timer.C:
             // 任務執行時間信號
             // to do something...
          case t := <-q.add:
             // 任務推送信號
             q.addTask(t)
          case id := <-q.remove:
             // 任務刪除信號
             // to do something...
          }
       }
    }
    6、生產者刪除任務
    // 用戶刪除任務
    func (q *DelayQueue) Delete(id string) {
       q.remove <- id
    }
    7、任務刪除信號的處理

    在這里,我們可以遍歷任務列表,根據刪除任務的id找到其在切片中的對應index。

    // 刪除指定任務
    func (q *DelayQueue) deleteTask(id string) {
       deleteIndex := -1
       for index, t := range q.tasks {
          if t.id == id {
             // 找到了在切片中需要刪除的所以呢
             deleteIndex = index
             break
          }
       }
       if deleteIndex == -1 {
          // 如果沒有找到刪除的任務,說明任務還在add管道中,來不及更新到tasks中,這里我們就將這個刪除id臨時記錄下來
          // 注意,這里暫時不考慮,任務id非法的特殊情況
          q.waitRemoveTaskMapping[id] = struct{}{}
          return
       }
       if len(q.tasks) == 1 {
          // 刪除后,任務列表就沒有任務了
          q.tasks = []*task{}
          return
       }
       if deleteIndex == len(q.tasks)-1 {
          // 如果刪除的是,任務列表的最后一個元素,則執行下列代碼
          q.tasks = q.tasks[:len(q.tasks)-1]
          return
       }
       // 如果刪除的是,任務列表的其他元素,則需要將deleteIndex之后的元素,全部向前挪動一位
       copy(q.tasks[deleteIndex:len(q.tasks)-1], q.tasks[deleteIndex+1:len(q.tasks)-1])
       q.tasks = q.tasks[:len(q.tasks)-1]
       return
    }

    然后,我們可以完善start()方法了。

    func (q *DelayQueue) start() {
       for {
          // to do something...
          select {
          case now := <-timer.C:
             // 任務執行時間信號
             // to do something...
          case t := <-q.add:
             // 任務推送信號
             q.addTask(t)
          case id := <-q.remove:
             // 任務刪除信號
             q.deleteTask(id)
          }
       }
    }
    8、任務執行信號的處理

    start()執行的時候,分成兩種情況:任務列表為空,只需要監聽add管道即可;任務列表不為空的時候,需要監聽所有管道。任務執行信號,主要是依靠timer來實現,屬于第二種情況。

    func (q *DelayQueue) start() {
       for {
          if len(q.tasks) == 0 {
               // 任務列表為空的時候,只需要監聽add管道
               select {
               case t := <-q.add:
                  //添加任務
                  q.addTask(t)
               }
               continue
          }
          // 任務列表不為空的時候,需要監聽所有管道
          // 任務的等待時間=任務的執行時間-當前的時間
          currentTask := q.tasks[0]
          timer := time.NewTimer(currentTask.execTime.Sub(time.Now()))
          select {
          case now := <-timer.C:
             // 任務執行信號
             timer.Stop()
            if _, isRemove := q.waitRemoveTaskMapping[currentTask.id]; isRemove {
               // 之前客戶已經發出過該任務的刪除信號,因此需要結束任務,刷新任務列表
               q.endTask()
               delete(q.waitRemoveTaskMapping, currentTask.id)
               continue
            }
            // 開啟協程,異步執行任務
            go q.execTask(currentTask, now)
            // 任務結束,刷新任務列表
            q.endTask()
          case t := <-q.add:
             // 任務推送信號
             timer.Stop()
             q.addTask(t)
          case id := <-q.remove:
             // 任務刪除信號
             timer.Stop()
             q.deleteTask(id)
          }
       }
    }

    執行任務:

    // 執行任務
    func (q *DelayQueue) execTask(task *task, currentTime time.Time) {
       if task.execTime.After(currentTime) {
          // 如果當前任務的執行時間落后于當前時間,則不執行
          return
       }
       // 執行任務
       task.f()
       return
    }

    結束任務,刷新任務列表:

    // 一個任務去執行了,刷新任務列表
    func (q *DelayQueue) endTask() {
       if len(q.tasks) == 1 {
          q.tasks = []*task{}
          return
       }
       q.tasks = q.tasks[1:]
    }
    9、完整代碼

    delay_queue.go

    package delay_queue
    import (
       "go.mongodb.org/mongo-driver/bson/primitive"
       "time"
    )
    // 延時任務對象
    type DelayQueue struct {
       tasks                 []*task             // 存儲任務列表的切片
       add                   chan *task          // 用戶添加任務的管道信號
       remove                chan string         // 用戶刪除任務的管道信號
       waitRemoveTaskMapping map[string]struct{} // 等待刪除的任務id列表
    }
    // 任務對象
    type task struct {
       id       string    // 任務id
       execTime time.Time // 執行時間
       f        func()    // 執行函數
    }
    // 初始化延時任務對象
    func InitDelayQueue() *DelayQueue {
       q := &DelayQueue{
          add:                   make(chan *task, 10000),
          remove:                make(chan string, 100),
          waitRemoveTaskMapping: make(map[string]struct{}),
       }
       // 開啟協程,監聽任務相關信號
       go q.start()
       return q
    }
    // 用戶刪除任務
    func (q *DelayQueue) Delete(id string) {
       q.remove <- id
    }
    // 用戶推送任務
    func (q *DelayQueue) Push(timeInterval time.Duration, f func()) string {
       // 生成一個任務id,方便刪除使用
       id := genTaskId()
       t := &task{
          id:       id,
          execTime: time.Now().Add(timeInterval),
          f:        f,
       }
       // 將任務推到add管道中
       q.add <- t
       return id
    }
    // 監聽各種任務相關信號
    func (q *DelayQueue) start() {
       for {
          if len(q.tasks) == 0 {
             // 任務列表為空的時候,只需要監聽add管道
             select {
             case t := <-q.add:
                //添加任務
                q.addTask(t)
             }
             continue
          }
          // 任務列表不為空的時候,需要監聽所有管道
          // 任務的等待時間=任務的執行時間-當前的時間
          currentTask := q.tasks[0]
          timer := time.NewTimer(currentTask.execTime.Sub(time.Now()))
          select {
          case now := <-timer.C:
             timer.Stop()
             if _, isRemove := q.waitRemoveTaskMapping[currentTask.id]; isRemove {
                // 之前客戶已經發出過該任務的刪除信號,因此需要結束任務,刷新任務列表
                q.endTask()
                delete(q.waitRemoveTaskMapping, currentTask.id)
                continue
             }
             // 開啟協程,異步執行任務
             go q.execTask(currentTask, now)
             // 任務結束,刷新任務列表
             q.endTask()
          case t := <-q.add:
             // 添加任務
             timer.Stop()
             q.addTask(t)
          case id := <-q.remove:
             // 刪除任務
             timer.Stop()
             q.deleteTask(id)
          }
       }
    }
    // 執行任務
    func (q *DelayQueue) execTask(task *task, currentTime time.Time) {
       if task.execTime.After(currentTime) {
          // 如果當前任務的執行時間落后于當前時間,則不執行
          return
       }
       // 執行任務
       task.f()
       return
    }
    // 一個任務去執行了,刷新任務列表
    func (q *DelayQueue) endTask() {
       if len(q.tasks) == 1 {
          q.tasks = []*task{}
          return
       }
       q.tasks = q.tasks[1:]
    }
    // 將任務添加到任務切片列表中
    func (q *DelayQueue) addTask(t *task) {
       // 尋找新增任務的插入位置
       insertIndex := q.getTaskInsertIndex(t, 0, len(q.tasks)-1)
       // 找到了插入位置,更新任務列表
       q.tasks = append(q.tasks, &task{})
       copy(q.tasks[insertIndex+1:], q.tasks[insertIndex:])
       q.tasks[insertIndex] = t
    }
    // 刪除指定任務
    func (q *DelayQueue) deleteTask(id string) {
       deleteIndex := -1
       for index, t := range q.tasks {
          if t.id == id {
             // 找到了在切片中需要刪除的所以呢
             deleteIndex = index
             break
          }
       }
       if deleteIndex == -1 {
          // 如果沒有找到刪除的任務,說明任務還在add管道中,來不及更新到tasks中,這里我們就將這個刪除id臨時記錄下來
          // 注意,這里暫時不考慮,任務id非法的特殊情況
          q.waitRemoveTaskMapping[id] = struct{}{}
          return
       }
       if len(q.tasks) == 1 {
          // 刪除后,任務列表就沒有任務了
          q.tasks = []*task{}
          return
       }
       if deleteIndex == len(q.tasks)-1 {
          // 如果刪除的是,任務列表的最后一個元素,則執行下列代碼
          q.tasks = q.tasks[:len(q.tasks)-1]
          return
       }
       // 如果刪除的是,任務列表的其他元素,則需要將deleteIndex之后的元素,全部向前挪動一位
       copy(q.tasks[deleteIndex:len(q.tasks)-1], q.tasks[deleteIndex+1:len(q.tasks)-1])
       q.tasks = q.tasks[:len(q.tasks)-1]
       return
    }
    // 尋找任務的插入位置
    func (q *DelayQueue) getTaskInsertIndex(t *task, leftIndex, rightIndex int) (index int) {
       // 使用二分法判斷新增任務的插入位置
       if len(q.tasks) == 0 {
          return
       }
       length := rightIndex - leftIndex
       if q.tasks[leftIndex].execTime.Sub(t.execTime) >= 0 {
          // 如果當前切片中最小的元素都超過了插入的優先級,則插入位置應該是最左邊
          return leftIndex
       }
       if q.tasks[rightIndex].execTime.Sub(t.execTime) <= 0 {
          // 如果當前切片中最大的元素都沒超過插入的優先級,則插入位置應該是最右邊
          return rightIndex + 1
       }
       if length == 1 && q.tasks[leftIndex].execTime.Before(t.execTime) && q.tasks[rightIndex].execTime.Sub(t.execTime) >= 0 {
          // 如果插入的優先級剛好在僅有的兩個優先級之間,則中間的位置就是插入位置
          return leftIndex + 1
       }
       middleVal := q.tasks[leftIndex+length/2].execTime
       // 這里用二分法遞歸的方式,一直尋找正確的插入位置
       if t.execTime.Sub(middleVal) <= 0 {
          return q.getTaskInsertIndex(t, leftIndex, leftIndex+length/2)
       } else {
          return q.getTaskInsertIndex(t, leftIndex+length/2, rightIndex)
       }
    }
    func genTaskId() string {
       return primitive.NewObjectID().Hex()
    }

    測試代碼:delay_queue_test.go

    package delay_queue
    import (
       "fmt"
       "testing"
       "time"
    )
    func TestDelayQueue(t *testing.T) {
       q := InitDelayQueue()
       for i := 0; i < 100; i++ {
          go func(i int) {
             id := q.Push(time.Duration(i)*time.Second, func() {
                fmt.Printf("%d秒后執行...\n", i)
                return
             })
             if i%7 == 0 {
                q.Delete(id)
             }
          }(i)
       }
       time.Sleep(time.Hour)
    }

    頭腦風暴

    上面的方案,的確實現了延時任務的效果,但是其中仍然有一些問題,仍然值得我們思考和優化。

    1、按照上面的方案,如果大量延時任務的執行時間,集中在同一個時間點,會造成短時間內timer頻繁地創建和銷毀。

    2、上述方案相比于time.AfterFunc()方法,我們需要在哪些場景下作出取舍。

    3、如果服務崩潰或重啟,如何去持久化隊列中的任務。

    到此,關于“golang延時任務如何實現”的學習就結束了,希望能夠解決大家的疑惑。理論與實踐的搭配能更好的幫助大家學習,快去試試吧!若想繼續學習更多相關知識,請繼續關注億速云網站,小編會繼續努力為大家帶來更多實用的文章!

    向AI問一下細節

    免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

    AI

    松滋市| 阳高县| 达日县| 永胜县| 寻甸| 泽库县| 永泰县| 鄱阳县| 饶河县| 汶上县| 明水县| 合水县| 新建县| 娱乐| 鄢陵县| 寿宁县| 兰坪| 濮阳市| 唐山市| 延津县| 尼勒克县| 云霄县| 大埔县| 石屏县| 漾濞| 棋牌| 孙吴县| 郧西县| 许昌县| 黎川县| 嘉善县| 高唐县| 阜宁县| 高青县| 盱眙县| 河东区| 乳源| 云安县| 淮北市| 罗甸县| 福泉市|