如何保证Go语言select中ticker分支不被通道读取抢占
Go批量更新任务的定时触发问题解决方案
问题背景
以下是两段用于批量处理用户积分更新的Go代码:
pointsQueue = make(chan *mongo.UpdateOneModel, 1000) func UpdatePoints(username string, size int64) { pointsDifference := -1 * size update := bson.D{{"$inc", bson.D{{"pointsLeft", pointsDifference }}}} updateOp := mongo.NewUpdateOneModel() updateOp.SetFilter(bson.M{"user": username}) updateOp.SetUpdate(update) pointsQueue <- updateOp } func updatePointsWorker() { var ctx = context.Background() ticker := time.NewTicker(dbBatchTimeout) var bulkRequests []mongo.WriteModel for { select { case req := <-pointsQueue: bulkRequests = append(bulkRequests, req) case <-ticker.C: if len(bulkRequests) > 0 { _, err := usersCollection.BulkWrite(ctx, bulkRequests) if err != nil { fmt.Println(err.Error()) } bulkRequests = nil } } }
当UpdatePoints在数秒内被调用数千次时,updatePointsWorker中的select会一直选中通道读取分支,导致ticker分支无法执行,批量请求无法定时清空。尝试过缓冲通道和批量阈值方案,但前者无效,后者会导致通道内的老请求因持续被新请求抢占而丢失。
解决方案
核心思路是避免通道读取完全阻塞定时任务的执行,同时保证所有请求都能被处理,以下是两种可靠的实现方式:
方式一:非阻塞读取+定时优先触发
修改worker逻辑,让定时事件优先得到处理,同时非阻塞地批量读取通道请求,避免霸占CPU:
func updatePointsWorker() { ctx := context.Background() ticker := time.NewTicker(dbBatchTimeout) defer ticker.Stop() // 回收ticker资源,避免泄露 var bulkRequests []mongo.WriteModel const batchSize = 500 // 自定义批量处理阈值 for { select { case <-ticker.C: // 定时触发批量写入 if len(bulkRequests) > 0 { if _, err := usersCollection.BulkWrite(ctx, bulkRequests); err != nil { fmt.Println(err.Error()) } bulkRequests = nil } default: // 非阻塞读取通道,每次最多读取batchSize个请求 readCount := 0 for readCount < batchSize { select { case req := <-pointsQueue: bulkRequests = append(bulkRequests, req) readCount++ // 达到批量阈值立即写入 if len(bulkRequests) >= batchSize { if _, err := usersCollection.BulkWrite(ctx, bulkRequests); err != nil { fmt.Println(err.Error()) } bulkRequests = nil readCount = 0 } default: // 通道无数据时退出内层循环 break } } // 让出CPU时间片,避免空转 time.Sleep(time.Millisecond * 1) } } }
逻辑说明:
- 外层
select优先处理ticker事件,确保定时任务不会被饿死 - 内层通过非阻塞读取,控制单次循环的读取数量,避免一直占用CPU
- 同时保留批量阈值,达到数量立即写入,兼顾处理效率和及时性
方式二:带超时的批量读取
在读取通道时加入超时机制,即使通道持续有数据,超时后也会回到外层select,给ticker执行机会:
func updatePointsWorker() { ctx := context.Background() ticker := time.NewTicker(dbBatchTimeout) defer ticker.Stop() var bulkRequests []mongo.WriteModel const batchSize = 500 const readTimeout = time.Millisecond * 10 // 单次读取超时时间,可按需调整 for { select { case <-ticker.C: // 定时写入逻辑不变 if len(bulkRequests) > 0 { if _, err := usersCollection.BulkWrite(ctx, bulkRequests); err != nil { fmt.Println(err.Error()) } bulkRequests = nil } case req := <-pointsQueue: bulkRequests = append(bulkRequests, req) // 达到批量阈值立即写入 if len(bulkRequests) >= batchSize { if _, err := usersCollection.BulkWrite(ctx, bulkRequests); err != nil { fmt.Println(err.Error()) } bulkRequests = nil } // 带超时读取更多请求,避免无限阻塞 readLoop: for len(bulkRequests) < batchSize { select { case req := <-pointsQueue: bulkRequests = append(bulkRequests, req) case <-time.After(readTimeout): break readLoop // 超时退出,回到外层select } } } } }
逻辑说明:
- 每次读取到一个请求后,进入带超时的内层循环,最多读满
batchSize个请求 - 超时后停止读取,回到外层
select,确保ticker有机会执行 - 超时时间可根据业务QPS调整,平衡处理效率和定时任务的及时性
关键注意事项
- 必须调用
ticker.Stop()回收资源,避免内存泄露 - 批量写入后将
bulkRequests重置为nil,而非重新创建切片,减少内存开销 - 通道缓冲大小可根据实际请求量调整,但核心是worker的读取逻辑,不能让通道读取完全阻塞定时任务
- 若担心批量写入失败,可添加重试机制,或把失败请求重新放回通道(需注意避免死循环)
内容的提问来源于stack exchange,提问作者Miriam Scapece
相关产品推荐
相关产品推荐

