并发编程
归纳 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,记录最终状态。