Skip to content
All notes

并发编程

归纳 Go 的 channel、context、同步原语和并发任务退出方式。

并发编程

channel

channel 兼顾传值和同步。底层 hchan 包含锁、可选的环形缓冲区,以及发送、接收等待队列。

字段作用
buf、dataqsiz、qcount缓冲区、容量、当前元素数
sendx、recvx下次写入、读取的位置,到末尾后绕回开头
sendq、recvq等待发送、接收的 goroutine 信息,不是待传的值列表
closed、lock关闭状态和内部同步

收发与关闭

状态发送接收
无缓冲等待接收方配对等待发送方配对
有缓冲、未关闭有空间就写入,满时等待有数据就读取,空且无发送者时等待
已关闭panic先读完缓冲数据,随后立即返回零值,ok=false
nil永久阻塞永久阻塞

close 表示不再发送,不会清空已有数据。重复关闭或关闭 nil channel 都会 panic;关闭责任应由能确定所有发送已结束的一方承担。

运行时的收发路径:

  • 有接收者等待时,发送者可以直接交付数据;否则先尝试写缓冲区,无法完成再进入 sendq。
  • 无缓冲 channel 有发送者等待时,接收者直接取值。
  • 有缓冲 channel 已满且有发送者等待时,先取缓冲区头部,再把等待发送者的值补到尾部,保持 FIFO;不会跳过旧数据。
  • 接收时既无数据又无发送者,则进入 recvq,挂起 goroutine,等待唤醒。

有缓冲 channel 的两个不变量:qcount > 0 时 recvq 为空;qcount < dataqsiz 时 sendq 为空。参见 channel 源码。

context

context 沿调用链传递取消信号、截止时间和少量请求级数据。父 context 取消,会传播给派生的子 context;取消只是通知,不会强制终止 goroutine,执行代码必须响应它。

方法作用
Done()取消或超时时关闭的只读 channel;永不取消的 context 可以返回 nil
Err()返回取消原因:Canceled 或 DeadlineExceeded;未取消时为 nil
Deadline()返回截止时间及是否设置
Value(key)读取请求级数据

Background 是根 context,TODO 用于尚未确定上下文的地方;WithCancel 支持手动取消,WithTimeout 和 WithDeadline 增加时间限制。拿到 cancel 后应及时调用,释放关联引用和计时器:

ctx, cancel := context.WithTimeout(req.Context(), 30*time.Second)
defer cancel()

把 ctx 传给下游 HTTP、RPC、数据库和模型请求。用户断连或超时后,相关工作才能及时停止。必须完成的持久化、计费或审计工作,需要单独定义完成责任和超时,不能一概随用户断连丢弃。

WithValue 用于 trace ID、请求身份等跨调用边界的数据,不传普通业务参数或大对象。参见 context 文档。

同步原语

原语用途边界
Mutex互斥访问共享状态同一时刻一个持有者;不绑定 goroutine
RWMutex多读者并行,写者独占等待中的写者会阻止新读者进入;不支持读锁升级,也不要递归加读锁
WaitGroup等待一组任务完成不保护共享数据,也不负责取消或传递错误

三者零值可用,首次使用后不能复制。一般先用 Mutex;读操作明显更多、读临界区有一定开销时,再测量 RWMutex 是否有收益。

读多写少时,锁与不可变快照如何取舍,可看 Pixiu 路由更新案例。

Mutex 的两种模式

正常模式下,被唤醒的等待者仍需与新来的 goroutine 竞争,吞吐量较好。源码中等待超过约 1ms 会触发饥饿模式:解锁者直接把锁交给队首等待者,减少长时间等待。

当接手的等待者已经是最后一个,或等待时间低于阈值时,锁回到正常模式。这是吞吐与公平性的取舍,阈值属于实现细节。

RWMutex 在互斥锁基础上结合读者计数和信号量协调读写,也有同步开销。参见 Mutex 源码 和 RWMutex 源码。

WaitGroup

var wg sync.WaitGroup
wg.Add(1)
go func() {
    defer wg.Done()
    // 执行任务
}()
wg.Wait()

Add 放在启动对应 goroutine 之前,避免 Wait 提前返回;Done 等价于 Add(-1),计数为负会 panic。复用时,下一批任务要等上一批 Wait 全部返回。

Go 1.25 起可用 wg.Go(f) 启动并计数任务,函数返回时自动扣减计数;文档要求 f 不得 panic。底层通过原子状态记录任务数和等待者,再用信号量阻塞、唤醒。参见 sync 文档。

并发模型:串行管理状态,限制并发执行

EventLoop 让一个 goroutine 持有状态,其他 goroutine 只投递事件。它适合会话状态机、流式输出聚合和任务编排。单消费者队列是其中一种形式:多个生产者投递,一个消费者按接收顺序处理。

下面的 jobs、process 代表业务任务和处理函数:

for {
    select {
    case <-ctx.Done():
        return
    case job, ok := <-jobs:
        if !ok {
            return
        }
        process(ctx, job)
    }
}

处理函数本身也要支持取消;慢 I/O 不宜阻塞管理主状态的循环。多个 channel 同时就绪时,select 不保证业务优先级或跨 channel 的先后顺序。

队列与退出

  • 队列设置容量;满时明确是等待、拒绝还是丢弃,避免无上限积压。
  • 发送和接收都要考虑取消,防止消费者退出后生产者永远阻塞。
  • 需要批处理时,按数量或时间触发提交;退出时明确剩余任务是处理完还是放弃。
  • 临时错误可以有限重试并退避;参数、权限等错误应直接返回。

工具调用链

一次 Agent 请求可以按 规划 → 调用工具 → 汇总结果 → 生成回复 → 结束 推进。EventLoop 管理状态,固定数量的 worker 执行工具,结果再返回 EventLoop;worker 不直接改主状态。

需要确定的边界:

  • 关联与去重:CallID 关联请求和结果;有副作用的操作还需要业务幂等机制,单靠 ID 不能防止重复执行。
  • 资源预算:限制队列容量、worker 数、调用步数和重试次数。先无限创建 goroutine 再等信号量,仍会积压 goroutine。
  • 超时:请求有总超时,每个工具也有自己的超时;请求相关工作共享取消信号。
  • 收尾:停止接收新任务,按约定处理剩余任务,取消并等待 worker,记录最终状态。