登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  Golang >  Go教程

Go sync.Cond 如何实现可唤醒的有界任务队列:处理满队列与空队列

来源:17golang原创

时间:2026-08-28 00:27:56 228浏览 收藏

后台任务突然把内存顶高时,第一眼往往只会盯着 goroutine 数量,却忽略了“生产速度超过消费速度”这条更直接的线索。给任务队列加一个固定容量,可以把这类压力挡在入口;再用 sync.Cond 让生产者在队列满时等待、消费者在队列空时等待,队列就不会靠轮询消耗 CPU。

有界队列的关键不是调用一次 Wait,而是让“队列未满”“队列非空”“队列已关闭”三个状态都在同一把锁下判断,并在每次唤醒后重新确认。

要点速览

  • 固定容量限制生产者占用的内存,满队列时等待而不是继续追加。
  • notFull 只负责唤醒入队方,notEmpty 只负责唤醒出队方。
  • Wait 必须放在条件循环里,不能用一次 if 判断代替。
  • 关闭时同时设置 closed 并调用 Broadcast,让所有等待者都有机会退出。

先把队列的生产约束说清楚

下面的队列只保存整数任务,容量固定为 2。生产者调用 Enqueue:当队列已经达到容量,就等待“未满”信号;消费者调用 Dequeue:当队列为空,就等待“非空”信号。这个例子不负责重试、持久化或任务超时,边界只放在并发状态本身。

这里用两组条件变量表达两个方向:notFull 对应“还能放入”,notEmpty 对应“可以取出”。它们共享同一个 sync.Mutex,因为条件判断和切片修改必须是一个不可拆开的动作。

用两个等待点实现最小队列

package main

import (
    "fmt"
    "sync"
)

type Queue struct {
    mu       sync.Mutex
    notFull  *sync.Cond
    notEmpty *sync.Cond
    items    []int
    capacity int
    closed   bool
}

func NewQueue(capacity int) *Queue {
    q := &Queue{capacity: capacity}
    q.notFull = sync.NewCond(&q.mu)
    q.notEmpty = sync.NewCond(&q.mu)
    return q
}

func (q *Queue) Enqueue(v int) bool {
    q.mu.Lock()
    defer q.mu.Unlock()
    for len(q.items) == q.capacity && !q.closed {
        q.notFull.Wait()
    }
    if q.closed {
        return false
    }
    q.items = append(q.items, v)
    q.notEmpty.Signal()
    return true
}

func (q *Queue) Dequeue() (int, bool) {
    q.mu.Lock()
    defer q.mu.Unlock()
    for len(q.items) == 0 && !q.closed {
        q.notEmpty.Wait()
    }
    if len(q.items) == 0 && q.closed {
        return 0, false
    }
    v := q.items[0]
    q.items = q.items[1:]
    q.notFull.Signal()
    return v, true
}

func (q *Queue) Close() {
    q.mu.Lock()
    q.closed = true
    q.mu.Unlock()
    q.notFull.Broadcast()
    q.notEmpty.Broadcast()
}

func main() {
    q := NewQueue(2)
    q.Enqueue(10)
    q.Enqueue(20)
    fmt.Println(q.Dequeue())
    q.Close()
}

入队成功后只需要唤醒一个消费者,所以使用 notEmpty.Signal();出队腾出空间后只需要唤醒一个生产者,所以使用 notFull.Signal()。如果队列关闭,可能同时存在两边的等待者,Close 必须对两组条件都调用 Broadcast

Go sync.Cond 有界队列中 Enqueue、notFull 与 Dequeue 的状态变化

为什么 Wait 一定要放在 for 里

Wait 返回只代表当前 goroutine 被唤醒,并不代表它重新获得锁后条件依然成立。多个消费者同时等待时,一个任务只够其中一个消费者取走;另一个消费者即使被唤醒,也要重新检查 len(q.items) == 0。生产者遇到满队列时同理,前一个生产者醒来后可能已经把最后一个空位占掉。

因此下面这种写法是不安全的:

if len(q.items) == q.capacity {
    q.notFull.Wait()
}
q.items = append(q.items, v)

安全写法是让循环同时检查条件和关闭状态:

for len(q.items) == q.capacity && !q.closed {
    q.notFull.Wait()
}
if q.closed {
    return false
}

这段代码的边界很具体:唤醒后先判断是否关闭,再判断是否真的有空间。不要把锁解开后再检查切片长度,否则判断结果和修改动作之间会被别的 goroutine 插入。

关闭状态要让等待者真正走出去

仅仅把 closed 改成 true 不够。正在 notFull.Wait()notEmpty.Wait() 的 goroutine 不会因为普通字段变化自动醒来,所以 Close 要先写入关闭状态,再广播两组条件。

出队方的语义是“关闭后把剩余任务取完,队列为空时返回 false”;入队方的语义是“关闭后不再接收新任务”。这两个判断不能合成一个模糊的返回值,否则调用方无法区分“暂时没任务”和“队列已经结束”。

Go sync.Cond 队列 Close 设置 closed 后 Broadcast 唤醒等待者并安全退出

并发验收:看三个状态而不是只看输出

把示例保存为 main.go 后先运行 go run -race main.go。如果把生产者和消费者改成 goroutine,再用 go test -race ./...,重点观察是否出现对 itemsclosed 的竞争。竞态检测通过,只能说明访问同步;它不能证明关闭语义符合业务预期。

  • 满队列:生产者是否停在 notFull.Wait(),出队后是否继续。
  • 空队列:消费者是否停在 notEmpty.Wait(),入队后是否拿到任务。
  • 关闭队列:等待中的两类 goroutine 是否都能被 Broadcast 唤醒并返回。

生产环境还要明确容量来源和关闭责任。容量过小会让生产者频繁等待,容量过大又会把突发压力推迟到内存;如果任务必须可靠保存,内存队列也不能替代持久化队列。

常见问题

sync.Cond 能不能只使用一个条件变量?

可以,但调用方需要在同一个条件变量上广播更多无关等待者,判断和唤醒关系不如 notFullnotEmpty 清晰。两个条件变量更容易审计每个状态转换。

Signal 和 Broadcast 应该怎么选?

单个任务到达或单个空位出现时通常用 Signal;关闭、配置整体变化或所有等待者都需要重新判断时用 Broadcast

为什么不直接用带缓冲 channel?

如果需求只是固定容量的生产消费,带缓冲 channel 通常更简单。需要显式关闭状态、批量唤醒或把多个条件纳入同一把锁时,sync.Cond 才更有发挥空间。

小结

这个队列真正受保护的是状态转换:入队让队列从“空”走向“非空”,出队让队列从“满”走向“未满”,关闭让所有等待者获得最后一次判断机会。把条件检查、切片修改和信号通知放在同一套锁协议里,再用 for 重新验证,才是 sync.Cond 处理有界任务队列的可靠基础。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>