溫馨提示×

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

密碼登錄×
登錄注冊(cè)×
其他方式登錄
點(diǎn)擊 登錄注冊(cè) 即表示同意《億速云用戶服務(wù)條款》

golang如何實(shí)現(xiàn)延時(shí)任務(wù)

發(fā)布時(shí)間:2023-03-23 09:46:23 來源:億速云 閱讀:182 作者:iii 欄目:編程語言

這篇文章主要講解了“golang如何實(shí)現(xiàn)延時(shí)任務(wù)”,文中的講解內(nèi)容簡單清晰,易于學(xué)習(xí)與理解,下面請(qǐng)大家跟著小編的思路慢慢深入,一起來研究和學(xué)習(xí)“golang如何實(shí)現(xiàn)延時(shí)任務(wù)”吧!

實(shí)現(xiàn)思路

我們都知道,任何一種隊(duì)列,實(shí)際上都是存在生產(chǎn)者和消費(fèi)者兩部分的。只不過,延時(shí)任務(wù)相對(duì)于普通隊(duì)列,多了一個(gè)延時(shí)的特性罷了。

1、生產(chǎn)者

從生產(chǎn)者的角度上講,當(dāng)用戶推送一個(gè)任務(wù)過來的時(shí)候,會(huì)攜帶著延遲執(zhí)行的時(shí)間數(shù)值。為了讓這個(gè)任務(wù)到預(yù)定時(shí)刻能執(zhí)行,我們需要將這個(gè)任務(wù)放在內(nèi)存里儲(chǔ)存一段時(shí)間,并且時(shí)間是一維的,在不斷增長。那么,我們用什么數(shù)據(jù)結(jié)構(gòu)存儲(chǔ)呢?

(1)選擇一:map。由于map具有無序性,無法按照?qǐng)?zhí)行時(shí)間排序,我們無法保證取出的任務(wù)是否是當(dāng)前時(shí)間點(diǎn)需要執(zhí)行的,所以排除這個(gè)選項(xiàng)。

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

(3)選擇三:slice。切片貌似可行,因?yàn)榍衅厥蔷哂杏行蛐缘?,所以,如果我們能夠按照?qǐng)?zhí)行時(shí)間的順序排列好所有的切片元素,那么,每次只要讀取切片的頭元素(也可能是尾元素),就可以得到我們要的任務(wù)。

2、消費(fèi)者

從消費(fèi)者的角度來說,它最大的難點(diǎn)在于,如何讓每個(gè)任務(wù),在特定的時(shí)間點(diǎn)被消費(fèi)。那么,針對(duì)每一個(gè)任務(wù),我們?nèi)绾螌?shí)現(xiàn),讓它等待一段時(shí)間后再執(zhí)行呢?

沒錯(cuò),就是timer。

總結(jié)下來,“切片+timer”的組合,應(yīng)該是可以達(dá)到目的的。

步步為營

1、數(shù)據(jù)流

(1)用戶調(diào)用InitDelayQueue() ,初始化延時(shí)任務(wù)對(duì)象。

(2)開啟協(xié)程,監(jiān)聽任務(wù)操作管道(add/delete信號(hào)),以及執(zhí)行時(shí)間管道(timer.C信號(hào))。

(3)用戶發(fā)出add/delete信號(hào)。

(4)(2)中的協(xié)程捕捉到(3)中的信號(hào),對(duì)任務(wù)列表進(jìn)行變更。

(5)當(dāng)任務(wù)執(zhí)行的時(shí)間點(diǎn)到達(dá)的時(shí)候(timer.C管道有元素輸出的時(shí)候),執(zhí)行任務(wù)。

golang如何實(shí)現(xiàn)延時(shí)任務(wù)

2、數(shù)據(jù)結(jié)構(gòu)

(1)延時(shí)任務(wù)對(duì)象

// 延時(shí)任務(wù)對(duì)象
type DelayQueue struct {
   tasks                 []*task             // 存儲(chǔ)任務(wù)列表的切片
   add                   chan *task          // 用戶添加任務(wù)的管道信號(hào)
   remove                chan string         // 用戶刪除任務(wù)的管道信號(hào)
   waitRemoveTaskMapping map[string]struct{} // 等待刪除的任務(wù)id列表
}

這里需要注意,有一個(gè)waitRemoveTaskMapping字段。由于要?jiǎng)h除的任務(wù),可能還在add管道中,沒有及時(shí)更新到tasks字段中,所以,需要臨時(shí)記錄下客戶要?jiǎng)h除的任務(wù)id。

(2)任務(wù)對(duì)象

// 任務(wù)對(duì)象
type task struct {
   id       string    // 任務(wù)id
   execTime time.Time // 執(zhí)行時(shí)間
   f        func()    // 執(zhí)行函數(shù)
}

3、初始化延時(shí)任務(wù)對(duì)象

// 初始化延時(shí)任務(wù)對(duì)象
func InitDelayQueue() *DelayQueue {
   q := &DelayQueue{
      add:                   make(chan *task, 10000),
      remove:                make(chan string, 100),
      waitRemoveTaskMapping: make(map[string]struct{}),
   }

   return q
}

在這個(gè)過程中,我們需要對(duì)用戶對(duì)任務(wù)的操作信號(hào),以及任務(wù)的執(zhí)行時(shí)間信號(hào)進(jìn)行監(jiān)聽。

func (q *DelayQueue) start() {
   for {
      // to do something...

      select {
      case now := <-timer.C:
         // 任務(wù)執(zhí)行時(shí)間信號(hào)
         // to do something...
      case t := <-q.add:
         // 任務(wù)推送信號(hào)
         // to do something...
      case id := <-q.remove:
         // 任務(wù)刪除信號(hào)
         // to do something...
      }
   }
}

完善我們的初始化方法:

// 初始化延時(shí)任務(wù)對(duì)象
func InitDelayQueue() *DelayQueue {
   q := &DelayQueue{
      add:                   make(chan *task, 10000),
      remove:                make(chan string, 100),
      waitRemoveTaskMapping: make(map[string]struct{}),
   }

   // 開啟協(xié)程,監(jiān)聽任務(wù)相關(guān)信號(hào)
   go q.start()
   return q
}

4、生產(chǎn)者推送任務(wù)

生產(chǎn)者推送任務(wù)的時(shí)候,只需要將任務(wù)加到add管道中即可,在這里,我們生成一個(gè)任務(wù)id,并返回給用戶。

// 用戶推送任務(wù)
func (q *DelayQueue) Push(timeInterval time.Duration, f func()) string {
   // 生成一個(gè)任務(wù)id,方便刪除使用
   id := genTaskId()
   t := &task{
      id:       id,
      execTime: time.Now().Add(timeInterval),
      f:        f,
   }

   // 將任務(wù)推到add管道中
   q.add <- t
   return id
}

5、任務(wù)推送信號(hào)的處理

在這里,我們要將用戶推送的任務(wù)放到延時(shí)任務(wù)的tasks字段中。由于,我們需要將任務(wù)按照?qǐng)?zhí)行時(shí)間順序排序,所以,我們需要找到新增任務(wù)在切片中的插入位置。又因?yàn)椋迦胫暗娜蝿?wù)列表已經(jīng)是有序的,所以,我們可以采用二分法處理。

// 使用二分法判斷新增任務(wù)的插入位置
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 {
      // 如果當(dāng)前切片中最小的元素都超過了插入的優(yōu)先級(jí),則插入位置應(yīng)該是最左邊
      return leftIndex
   }

   if q.tasks[rightIndex].execTime.Sub(t.execTime) <= 0 {
      // 如果當(dāng)前切片中最大的元素都沒超過插入的優(yōu)先級(jí),則插入位置應(yīng)該是最右邊
      return rightIndex + 1
   }

   if length == 1 && q.tasks[leftIndex].execTime.Before(t.execTime) && q.tasks[rightIndex].execTime.Sub(t.execTime) >= 0 {
      // 如果插入的優(yōu)先級(jí)剛好在僅有的兩個(gè)優(yōu)先級(jí)之間,則中間的位置就是插入位置
      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)
   }
}

找到正確的插入位置后,我們才能將任務(wù)準(zhǔn)確插入:

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

那么,在監(jiān)聽add管道的時(shí)候,我們直接調(diào)用上述addTask() 即可。

func (q *DelayQueue) start() {
   for {
      // to do something...

      select {
      case now := <-timer.C:
         // 任務(wù)執(zhí)行時(shí)間信號(hào)
         // to do something...
      case t := <-q.add:
         // 任務(wù)推送信號(hào)
         q.addTask(t)
      case id := <-q.remove:
         // 任務(wù)刪除信號(hào)
         // to do something...
      }
   }
}

6、生產(chǎn)者刪除任務(wù)

// 用戶刪除任務(wù)
func (q *DelayQueue) Delete(id string) {
   q.remove <- id
}

7、任務(wù)刪除信號(hào)的處理

在這里,我們可以遍歷任務(wù)列表,根據(jù)刪除任務(wù)的id找到其在切片中的對(duì)應(yīng)index。

// 刪除指定任務(wù)
func (q *DelayQueue) deleteTask(id string) {
   deleteIndex := -1
   for index, t := range q.tasks {
      if t.id == id {
         // 找到了在切片中需要?jiǎng)h除的所以呢
         deleteIndex = index
         break
      }
   }

   if deleteIndex == -1 {
      // 如果沒有找到刪除的任務(wù),說明任務(wù)還在add管道中,來不及更新到tasks中,這里我們就將這個(gè)刪除id臨時(shí)記錄下來
      // 注意,這里暫時(shí)不考慮,任務(wù)id非法的特殊情況
      q.waitRemoveTaskMapping[id] = struct{}{}
      return
   }

   if len(q.tasks) == 1 {
      // 刪除后,任務(wù)列表就沒有任務(wù)了
      q.tasks = []*task{}
      return
   }

   if deleteIndex == len(q.tasks)-1 {
      // 如果刪除的是,任務(wù)列表的最后一個(gè)元素,則執(zhí)行下列代碼
      q.tasks = q.tasks[:len(q.tasks)-1]
      return
   }

   // 如果刪除的是,任務(wù)列表的其他元素,則需要將deleteIndex之后的元素,全部向前挪動(dòng)一位
   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:
         // 任務(wù)執(zhí)行時(shí)間信號(hào)
         // to do something...
      case t := <-q.add:
         // 任務(wù)推送信號(hào)
         q.addTask(t)
      case id := <-q.remove:
         // 任務(wù)刪除信號(hào)
         q.deleteTask(id)
      }
   }
}

8、任務(wù)執(zhí)行信號(hào)的處理

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

func (q *DelayQueue) start() {
   for {
      if len(q.tasks) == 0 {
           // 任務(wù)列表為空的時(shí)候,只需要監(jiān)聽add管道
           select {
           case t := <-q.add:
              //添加任務(wù)
              q.addTask(t)
           }
        
           continue
      }

      // 任務(wù)列表不為空的時(shí)候,需要監(jiān)聽所有管道

      // 任務(wù)的等待時(shí)間=任務(wù)的執(zhí)行時(shí)間-當(dāng)前的時(shí)間
      currentTask := q.tasks[0]
      timer := time.NewTimer(currentTask.execTime.Sub(time.Now()))

      select {
      case now := <-timer.C:
         // 任務(wù)執(zhí)行信號(hào)
         timer.Stop()
        if _, isRemove := q.waitRemoveTaskMapping[currentTask.id]; isRemove {
           // 之前客戶已經(jīng)發(fā)出過該任務(wù)的刪除信號(hào),因此需要結(jié)束任務(wù),刷新任務(wù)列表
           q.endTask()
           delete(q.waitRemoveTaskMapping, currentTask.id)
           continue
        }
        
        // 開啟協(xié)程,異步執(zhí)行任務(wù)
        go q.execTask(currentTask, now)
        // 任務(wù)結(jié)束,刷新任務(wù)列表
        q.endTask()
      case t := <-q.add:
         // 任務(wù)推送信號(hào)
         timer.Stop()
         q.addTask(t)
      case id := <-q.remove:
         // 任務(wù)刪除信號(hào)
         timer.Stop()
         q.deleteTask(id)
      }
   }
}

執(zhí)行任務(wù):

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

   // 執(zhí)行任務(wù)
   task.f()
   return
}

結(jié)束任務(wù),刷新任務(wù)列表:

// 一個(gè)任務(wù)去執(zhí)行了,刷新任務(wù)列表
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"
)

// 延時(shí)任務(wù)對(duì)象
type DelayQueue struct {
   tasks                 []*task             // 存儲(chǔ)任務(wù)列表的切片
   add                   chan *task          // 用戶添加任務(wù)的管道信號(hào)
   remove                chan string         // 用戶刪除任務(wù)的管道信號(hào)
   waitRemoveTaskMapping map[string]struct{} // 等待刪除的任務(wù)id列表
}

// 任務(wù)對(duì)象
type task struct {
   id       string    // 任務(wù)id
   execTime time.Time // 執(zhí)行時(shí)間
   f        func()    // 執(zhí)行函數(shù)
}

// 初始化延時(shí)任務(wù)對(duì)象
func InitDelayQueue() *DelayQueue {
   q := &DelayQueue{
      add:                   make(chan *task, 10000),
      remove:                make(chan string, 100),
      waitRemoveTaskMapping: make(map[string]struct{}),
   }

   // 開啟協(xié)程,監(jiān)聽任務(wù)相關(guān)信號(hào)
   go q.start()
   return q
}

// 用戶刪除任務(wù)
func (q *DelayQueue) Delete(id string) {
   q.remove <- id
}

// 用戶推送任務(wù)
func (q *DelayQueue) Push(timeInterval time.Duration, f func()) string {
   // 生成一個(gè)任務(wù)id,方便刪除使用
   id := genTaskId()
   t := &task{
      id:       id,
      execTime: time.Now().Add(timeInterval),
      f:        f,
   }

   // 將任務(wù)推到add管道中
   q.add <- t
   return id
}

// 監(jiān)聽各種任務(wù)相關(guān)信號(hào)
func (q *DelayQueue) start() {
   for {
      if len(q.tasks) == 0 {
         // 任務(wù)列表為空的時(shí)候,只需要監(jiān)聽add管道
         select {
         case t := <-q.add:
            //添加任務(wù)
            q.addTask(t)
         }

         continue
      }

      // 任務(wù)列表不為空的時(shí)候,需要監(jiān)聽所有管道

      // 任務(wù)的等待時(shí)間=任務(wù)的執(zhí)行時(shí)間-當(dāng)前的時(shí)間
      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 {
            // 之前客戶已經(jīng)發(fā)出過該任務(wù)的刪除信號(hào),因此需要結(jié)束任務(wù),刷新任務(wù)列表
            q.endTask()
            delete(q.waitRemoveTaskMapping, currentTask.id)
            continue
         }

         // 開啟協(xié)程,異步執(zhí)行任務(wù)
         go q.execTask(currentTask, now)
         // 任務(wù)結(jié)束,刷新任務(wù)列表
         q.endTask()
      case t := <-q.add:
         // 添加任務(wù)
         timer.Stop()
         q.addTask(t)
      case id := <-q.remove:
         // 刪除任務(wù)
         timer.Stop()
         q.deleteTask(id)
      }
   }
}

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

   // 執(zhí)行任務(wù)
   task.f()
   return
}

// 一個(gè)任務(wù)去執(zhí)行了,刷新任務(wù)列表
func (q *DelayQueue) endTask() {
   if len(q.tasks) == 1 {
      q.tasks = []*task{}
      return
   }

   q.tasks = q.tasks[1:]
}

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

// 刪除指定任務(wù)
func (q *DelayQueue) deleteTask(id string) {
   deleteIndex := -1
   for index, t := range q.tasks {
      if t.id == id {
         // 找到了在切片中需要?jiǎng)h除的所以呢
         deleteIndex = index
         break
      }
   }

   if deleteIndex == -1 {
      // 如果沒有找到刪除的任務(wù),說明任務(wù)還在add管道中,來不及更新到tasks中,這里我們就將這個(gè)刪除id臨時(shí)記錄下來
      // 注意,這里暫時(shí)不考慮,任務(wù)id非法的特殊情況
      q.waitRemoveTaskMapping[id] = struct{}{}
      return
   }

   if len(q.tasks) == 1 {
      // 刪除后,任務(wù)列表就沒有任務(wù)了
      q.tasks = []*task{}
      return
   }

   if deleteIndex == len(q.tasks)-1 {
      // 如果刪除的是,任務(wù)列表的最后一個(gè)元素,則執(zhí)行下列代碼
      q.tasks = q.tasks[:len(q.tasks)-1]
      return
   }

   // 如果刪除的是,任務(wù)列表的其他元素,則需要將deleteIndex之后的元素,全部向前挪動(dòng)一位
   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
}

// 尋找任務(wù)的插入位置
func (q *DelayQueue) getTaskInsertIndex(t *task, leftIndex, rightIndex int) (index int) {
   // 使用二分法判斷新增任務(wù)的插入位置
   if len(q.tasks) == 0 {
      return
   }

   length := rightIndex - leftIndex
   if q.tasks[leftIndex].execTime.Sub(t.execTime) >= 0 {
      // 如果當(dāng)前切片中最小的元素都超過了插入的優(yōu)先級(jí),則插入位置應(yīng)該是最左邊
      return leftIndex
   }

   if q.tasks[rightIndex].execTime.Sub(t.execTime) <= 0 {
      // 如果當(dāng)前切片中最大的元素都沒超過插入的優(yōu)先級(jí),則插入位置應(yīng)該是最右邊
      return rightIndex + 1
   }

   if length == 1 && q.tasks[leftIndex].execTime.Before(t.execTime) && q.tasks[rightIndex].execTime.Sub(t.execTime) >= 0 {
      // 如果插入的優(yōu)先級(jí)剛好在僅有的兩個(gè)優(yōu)先級(jí)之間,則中間的位置就是插入位置
      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()
}

測(cè)試代碼: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秒后執(zhí)行...\n", i)
            return
         })

         if i%7 == 0 {
            q.Delete(id)
         }
      }(i)
   }

   time.Sleep(time.Hour)
}

頭腦風(fēng)暴

上面的方案,的確實(shí)現(xiàn)了延時(shí)任務(wù)的效果,但是其中仍然有一些問題,仍然值得我們思考和優(yōu)化。

1、按照上面的方案,如果大量延時(shí)任務(wù)的執(zhí)行時(shí)間,集中在同一個(gè)時(shí)間點(diǎn),會(huì)造成短時(shí)間內(nèi)timer頻繁地創(chuàng)建和銷毀。

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

3、如果服務(wù)崩潰或重啟,如何去持久化隊(duì)列中的任務(wù)。

感謝各位的閱讀,以上就是“golang如何實(shí)現(xiàn)延時(shí)任務(wù)”的內(nèi)容了,經(jīng)過本文的學(xué)習(xí)后,相信大家對(duì)golang如何實(shí)現(xiàn)延時(shí)任務(wù)這一問題有了更深刻的體會(huì),具體使用情況還需要大家實(shí)踐驗(yàn)證。這里是億速云,小編將為大家推送更多相關(guān)知識(shí)點(diǎn)的文章,歡迎關(guān)注!

向AI問一下細(xì)節(jié)

免責(zé)聲明:本站發(fā)布的內(nèi)容(圖片、視頻和文字)以原創(chuàng)、轉(zhuǎn)載和分享為主,文章觀點(diǎn)不代表本網(wǎng)站立場(chǎng),如果涉及侵權(quán)請(qǐng)聯(lián)系站長郵箱:is@yisu.com進(jìn)行舉報(bào),并提供相關(guān)證據(jù),一經(jīng)查實(shí),將立刻刪除涉嫌侵權(quán)內(nèi)容。

AI