一文詳解golang延時(shí)任務(wù)的實(shí)現(xiàn)
前言
在實(shí)際業(yè)務(wù)場景中,我們有時(shí)候會碰到一些延時(shí)的需求:例如,在電商平臺,運(yùn)營在管理后臺添加商品后,不需要立刻展示在前臺,而是在之后某個(gè)時(shí)間點(diǎn)才展現(xiàn)。
當(dāng)然,我們有很多種思路,可以應(yīng)對這個(gè)問題。例如,將待發(fā)布商品信息添加到db,然后通過定時(shí)任務(wù)輪詢數(shù)據(jù)表的方式,查詢當(dāng)前時(shí)間點(diǎn)的發(fā)布商品;又比如,將商品信息全部添加到redis中,通過SortSet屬性完成這個(gè)功能。最終的選擇,取決于我們的業(yè)務(wù)場景和運(yùn)行環(huán)境。
在這里,我想給大家分享一套,基于golang實(shí)現(xiàn)的延時(shí)任務(wù)方案。
你可以收獲
- golang管道的靈活運(yùn)用
- golang timer的應(yīng)用
- golang切片元素插入排序的實(shí)現(xiàn)思路
- golang延時(shí)任務(wù)的實(shí)現(xiàn)思路
正文
思維導(dǎo)圖
為了讓大家有一個(gè)大致的印象,我將正文的大綱列在下面。

實(shí)現(xiàn)思路
我們都知道,任何一種隊(duì)列,實(shí)際上都是存在生產(chǎn)者和消費(fèi)者兩部分的。只不過,延時(shí)任務(wù)相對于普通隊(duì)列,多了一個(gè)延時(shí)的特性罷了。
1、生產(chǎn)者
從生產(chǎn)者的角度上講,當(dāng)用戶推送一個(gè)任務(wù)過來的時(shí)候,會攜帶著延遲執(zhí)行的時(shí)間數(shù)值。為了讓這個(gè)任務(wù)到預(yù)定時(shí)刻能執(zhí)行,我們需要將這個(gè)任務(wù)放在內(nèi)存里儲存一段時(shí)間,并且時(shí)間是一維的,在不斷增長。那么,我們用什么數(shù)據(jù)結(jié)構(gòu)存儲呢?
(1)選擇一:map。由于map具有無序性,無法按照執(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)榍衅厥蔷哂杏行蛐缘模?,如果我們能夠按照?zhí)行時(shí)間的順序排列好所有的切片元素,那么,每次只要讀取切片的頭元素(也可能是尾元素),就可以得到我們要的任務(wù)。
2、消費(fèi)者
從消費(fèi)者的角度來說,它最大的難點(diǎn)在于,如何讓每個(gè)任務(wù),在特定的時(shí)間點(diǎn)被消費(fèi)。那么,針對每一個(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ù)對象。
(2)開啟協(xié)程,監(jiān)聽任務(wù)操作管道(add/delete信號),以及執(zhí)行時(shí)間管道(timer.C信號)。
(3)用戶發(fā)出add/delete信號。
(4)(2)中的協(xié)程捕捉到(3)中的信號,對任務(wù)列表進(jìn)行變更。
(5)當(dāng)任務(wù)執(zhí)行的時(shí)間點(diǎn)到達(dá)的時(shí)候(timer.C管道有元素輸出的時(shí)候),執(zhí)行任務(wù)。

2、數(shù)據(jù)結(jié)構(gòu)
(1)延時(shí)任務(wù)對象
// 延時(shí)任務(wù)對象
type DelayQueue struct {
tasks []*task // 存儲任務(wù)列表的切片
add chan *task // 用戶添加任務(wù)的管道信號
remove chan string // 用戶刪除任務(wù)的管道信號
waitRemoveTaskMapping map[string]struct{} // 等待刪除的任務(wù)id列表
}
這里需要注意,有一個(gè)waitRemoveTaskMapping字段。由于要?jiǎng)h除的任務(wù),可能還在add管道中,沒有及時(shí)更新到tasks字段中,所以,需要臨時(shí)記錄下客戶要?jiǎng)h除的任務(wù)id。
(2)任務(wù)對象
// 任務(wù)對象
type task struct {
id string // 任務(wù)id
execTime time.Time // 執(zhí)行時(shí)間
f func() // 執(zhí)行函數(shù)
}
3、初始化延時(shí)任務(wù)對象
// 初始化延時(shí)任務(wù)對象
func InitDelayQueue() *DelayQueue {
q := &DelayQueue{
add: make(chan *task, 10000),
remove: make(chan string, 100),
waitRemoveTaskMapping: make(map[string]struct{}),
}
return q
}
在這個(gè)過程中,我們需要對用戶對任務(wù)的操作信號,以及任務(wù)的執(zhí)行時(shí)間信號進(jìn)行監(jiān)聽。
func (q *DelayQueue) start() {
for {
// to do something...
select {
case now := <-timer.C:
// 任務(wù)執(zhí)行時(shí)間信號
// to do something...
case t := <-q.add:
// 任務(wù)推送信號
// to do something...
case id := <-q.remove:
// 任務(wù)刪除信號
// to do something...
}
}
}
完善我們的初始化方法:
// 初始化延時(shí)任務(wù)對象
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)信號
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ù)推送信號的處理
在這里,我們要將用戶推送的任務(wù)放到延時(shí)任務(wù)的tasks字段中。由于,我們需要將任務(wù)按照執(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)先級,則插入位置應(yīng)該是最左邊
return leftIndex
}
if q.tasks[rightIndex].execTime.Sub(t.execTime) <= 0 {
// 如果當(dāng)前切片中最大的元素都沒超過插入的優(yōu)先級,則插入位置應(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)先級剛好在僅有的兩個(gè)優(yōu)先級之間,則中間的位置就是插入位置
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í)間信號
// to do something...
case t := <-q.add:
// 任務(wù)推送信號
q.addTask(t)
case id := <-q.remove:
// 任務(wù)刪除信號
// to do something...
}
}
}
6、生產(chǎn)者刪除任務(wù)
// 用戶刪除任務(wù)
func (q *DelayQueue) Delete(id string) {
q.remove <- id
}
7、任務(wù)刪除信號的處理
在這里,我們可以遍歷任務(wù)列表,根據(jù)刪除任務(wù)的id找到其在切片中的對應(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í)間信號
// to do something...
case t := <-q.add:
// 任務(wù)推送信號
q.addTask(t)
case id := <-q.remove:
// 任務(wù)刪除信號
q.deleteTask(id)
}
}
}
8、任務(wù)執(zhí)行信號的處理
start()執(zhí)行的時(shí)候,分成兩種情況:任務(wù)列表為空,只需要監(jiān)聽add管道即可;任務(wù)列表不為空的時(shí)候,需要監(jiān)聽所有管道。任務(wù)執(zhí)行信號,主要是依靠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í)行信號
timer.Stop()
if _, isRemove := q.waitRemoveTaskMapping[currentTask.id]; isRemove {
// 之前客戶已經(jīng)發(fā)出過該任務(wù)的刪除信號,因此需要結(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ù):
// 執(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ù)對象
type DelayQueue struct {
tasks []*task // 存儲任務(wù)列表的切片
add chan *task // 用戶添加任務(wù)的管道信號
remove chan string // 用戶刪除任務(wù)的管道信號
waitRemoveTaskMapping map[string]struct{} // 等待刪除的任務(wù)id列表
}
// 任務(wù)對象
type task struct {
id string // 任務(wù)id
execTime time.Time // 執(zhí)行時(shí)間
f func() // 執(zhí)行函數(shù)
}
// 初始化延時(shí)任務(wù)對象
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)信號
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)信號
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ù)的刪除信號,因此需要結(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)先級,則插入位置應(yīng)該是最左邊
return leftIndex
}
if q.tasks[rightIndex].execTime.Sub(t.execTime) <= 0 {
// 如果當(dāng)前切片中最大的元素都沒超過插入的優(yōu)先級,則插入位置應(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)先級剛好在僅有的兩個(gè)優(yōu)先級之間,則中間的位置就是插入位置
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秒后執(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),會造成短時(shí)間內(nèi)timer頻繁地創(chuàng)建和銷毀。
2、上述方案相比于time.AfterFunc()方法,我們需要在哪些場景下作出取舍。
3、如果服務(wù)崩潰或重啟,如何去持久化隊(duì)列中的任務(wù)。
小結(jié)
本文和大家討論了延時(shí)任務(wù)在golang中的一種實(shí)現(xiàn)方案,在這個(gè)過程中,一次性定時(shí)器timer、切片、管道等golang特色,以及二分插入等常見算法都體現(xiàn)得淋漓盡致。
相關(guān)文章
Go語言中進(jìn)行API限流的實(shí)戰(zhàn)詳解
API?限流是控制和管理應(yīng)用程序訪問量的重要手段,旨在防止惡意濫用、保護(hù)后端服務(wù)的穩(wěn)定性和可用性,下面我們就來看看如何在Go語言中具體實(shí)現(xiàn)吧2025-01-01
Gin與Mysql實(shí)現(xiàn)簡單Restful風(fēng)格API實(shí)戰(zhàn)示例詳解
這篇文章主要為大家介紹了Gin與Mysql實(shí)現(xiàn)簡單Restful風(fēng)格API示例詳解,有需要的朋友可以借鑒參考下希望能夠有所幫助,祝大家多多進(jìn)步2021-11-11
Golang自動(dòng)追蹤GitHub上熱門AI項(xiàng)目
這篇文章主要為大家介紹了Golang自動(dòng)追蹤GitHub上熱門AI項(xiàng)目,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-12-12
golang?sync.Cond同步機(jī)制運(yùn)用及實(shí)現(xiàn)
在?Go?里有專門為同步通信而生的?channel,所以較少看到?sync.Cond?的使用,不過它也是并發(fā)控制手段里的一種,今天我們就來認(rèn)識下它的相關(guān)實(shí)現(xiàn),加深對同步機(jī)制的運(yùn)用2023-09-09
Go語言Gin框架中使用MySQL數(shù)據(jù)庫的三種方式
本文主要介紹了Go語言Gin框架中使用MySQL數(shù)據(jù)庫的三種方式,通過三種方式實(shí)現(xiàn)增刪改查的操作,具有一定的參考價(jià)值,感興趣的可以了解一下2023-11-11
go中的參數(shù)傳遞是值傳遞還是引用傳遞的實(shí)現(xiàn)
參數(shù)傳遞機(jī)制是一個(gè)重要的概念,它決定了函數(shù)內(nèi)部對參數(shù)的修改是否會影響到原始數(shù)據(jù),本文主要介紹了go中的參數(shù)傳遞是值傳遞還是引用傳遞的實(shí)現(xiàn),感興趣的可以了解一下2024-12-12

