You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何保证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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 08:30:00