Go 并发模式概览
常见进程内并发结构:pipeline、fan-out/fan-in、超时与取消、or-done、防止泄漏;与 worker pool 专题的分工。
[!info] 关联笔记
Go 并发模式概览
这个概念为什么会出现
掌握 goroutine、channel、select 之后,工程里反复撞上同一类结构问题:
- 多阶段处理如何串成流水线?
- 如何扇出并行又限制并发度?
- 多个结果如何汇合?
- 超时/取消时如何不泄漏 goroutine?
- 谁关闭 channel?完成信号如何传播?
这些问题不是再背一个 API 能解决的,而是可复用的进程内并发拓扑。本篇给模式地图与判定标准;具体落地可下钻 go-worker-pool、go-pipeline-pattern。
官方与社区长期材料包括 Go blog 的 concurrency patterns 系列、golang.org/x/sync 等。模式是工具箱,不是宗教——正确性(不泄漏、不 race、可取消)优先于“看起来很并发”。
[!abstract] 一句话理解 用 channel 表达数据流与完成信号,用 context 表达取消,用有界并发保护资源;每个
go都有明确退出路径与所有权规则。
最小可运行示例
先把示例放进业务场景,再看代码:
场景:订单 ID 流水线平方变换 + 支付回调超时
批处理要把一批订单号送进流水线:gen 产出 ID,sq 做“阶段变换”(这里用平方代替计价/归一化)。
发送方关闭 channel,下游 range 自然结束——这是 pipeline 的完成信号。
另一侧:等支付回调最多 10ms;下游 50ms 才回,必须走超时分支,
避免网关线程一直挂着。
package main
import (
"fmt"
"time"
)
// genOrderIDs 模拟“从批任务读出订单 ID 流”。
//
// 业务意图:只负责产出,产出完毕关闭 out,通知下游没有更多单。
// 教学点:发送方拥有 out 并 defer close;返回只读 channel 防止外人误关。
func genOrderIDs(ids ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, id := range ids {
out <- id
}
}()
return out
}
// squareStage 模拟流水线下一阶段(计价/哈希等纯变换)。
//
// 业务意图:读完上游就关闭自己的 out,把“完成”继续向后传。
// 教学点:阶段内拥有 out;上游关闭 → range 结束 → close(out)。
func squareStage(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
out <- n * n
}
}()
return out
}
// waitPaymentCallback 模拟“等支付渠道回调,但只愿等 budget”。
//
// 业务意图:回调慢就超时返回,释放请求槽;不能无限 block。
// 教学点:select + time.After(生产更常用 context 截止时间)。
func waitPaymentCallback(budget time.Duration) string {
ch := make(chan string, 1)
go func() {
// 假装渠道 50ms 后才回调。
time.Sleep(50 * time.Millisecond)
ch <- "paid"
}()
select {
case s := <-ch:
return s
case <-time.After(budget):
return "timeout"
}
}
func main() {
// --- 场景 A:订单 ID 两阶段 pipeline ---
for v := range squareStage(genOrderIDs(1, 2, 3)) {
fmt.Println(v)
}
// --- 场景 B:支付回调超时保护 ---
fmt.Println(waitPaymentCallback(10 * time.Millisecond))
}
建议运行:
go run .
期望输出:
1
4
9
timeout
结合场景再看三个关注点
-
发送方关闭,接收方 range
订单流结束靠close,不是另传“EOF 业务消息”(除非协议需要)。 -
阶段函数签名暴露只读 in/out
所有权清晰,避免“谁都可以 close”的竞态。 -
超时是并发模式的一部分
支付回调用select限制等待;可取消长链应再叠context。
核心概念与准确模型
1. 模式地图
| 模式 | 问题 | 关键构件 |
|---|---|---|
| Pipeline | 多阶段变换 | 阶段函数、close(out) |
| Fan-out | 提高吞吐 | 多 worker 读输入 |
| Fan-in | 合并多路 | WaitGroup + 单 closer |
| 有界并行 | 保护资源 | semaphore channel / 池 |
| 超时 | 限制等待 | select + timer/context |
| 取消 | 全链退出 | context |
| or-done | 释放阻塞 | 额外 done 通道 |
| 错误传播 | 多任务失败 | errgroup / err ch |
| Bridge / tee | 分发流 | 多输出或桥接 |
2. Pipeline(流水线)
每一阶段:
in <-chan T → (goroutine) → out <-chan U
defer close(out)
原则:
- 阶段内拥有
out的发送与关闭权 - 上游关闭 → 本阶段
range结束 → 关闭下游 - 数据与完成信号都走 channel
详见 go-pipeline-pattern。官方范例见 Go Blog — Pipelines。
3. Fan-out / Fan-in
Fan-out:多个 worker 从同一输入取任务(竞争消费)或从分片输入取。
Fan-in:多路输入合并到一个输出。
func merge(cs ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
wg.Add(len(cs))
for _, c := range cs {
c := c
go func() {
defer wg.Done()
for v := range c {
out <- v
}
}()
}
go func() {
wg.Wait()
close(out) // 唯一 closer
}()
return out
}
关键:
- 只有一个地方
close(out) - 用 WaitGroup 等所有发送者结束
- 见 go-channel-ownership-and-closing
4. 有界并行(semaphore)
无限 go func 可打爆内存/FD/下游。
sem := make(chan struct{}, maxConcurrent)
for _, item := range items {
sem <- struct{}{} // acquire
go func(it Item) {
defer func() { <-sem }() // release
handle(it)
}(item)
}
或固定 worker 数读任务队列——即 worker pool,见 go-worker-pool。
golang.org/x/sync/semaphore 提供加权信号量:semaphore。
5. 超时与取消
select {
case r := <-res:
// 使用 r
case <-ctx.Done():
return ctx.Err()
case <-time.After(dt):
return errTimeout
}
优先级建议:
- context 贯穿全链(请求级取消)
- 循环内慎用
time.After(每次分配 timer,历史版本易泄漏等待中的 timer;可用time.NewTimer并Stop) - 超时后仍要考虑后台 goroutine 是否退出
6. 防止 goroutine 泄漏
泄漏典型形态:
- 发送到无人接收且无缓冲 channel
range未关闭的 channel- 等待
ctx但从不知道 cancel - select 缺少 done 分支
检查清单:
- 每个
go的退出条件是什么? - 取消时阻塞在 channel 上的双方能否解锁?
- 测试里用
goleak或超时失败是否够?
7. or-done / 释放阻塞发送
当下游可能提前退出,上游发送会卡死:
func orDone(done <-chan struct{}, in <-chan T) <-chan T {
out := make(chan T)
go func() {
defer close(out)
for {
select {
case <-done:
return
case v, ok := <-in:
if !ok {
return
}
select {
case out <- v:
case <-done:
return
}
}
}
}()
return out
}
现代代码更多把 done 换成 ctx.Done()。
8. 错误传播模式
单错误通道:
errc := make(chan error, 1)
errgroup:
g, ctx := errgroup.WithContext(ctx)
g.Go(func() error { /* ... */ return nil })
err := g.Wait()
文档:errgroup
选择:
- 需要 cancel 兄弟任务 → errgroup.WithContext
- 收集所有错误 → 自己设计
[]error或errors.Join
9. 请求级并发 vs 进程级 worker
| 请求内 fan-out | 全局 worker pool | |
|---|---|---|
| 生命周期 | 随请求结束 | 随进程 |
| 取消 | 天然绑 request ctx | 任务自带 ctx |
| 资源 | 易打爆若无界 | 全局有界 |
API 服务常:全局有界 + 每任务 ctx。
10. 同步原语该上场时
channel 不是锤子唯一形状:
- 共享状态内存更新 →
sync.Mutex - 一次性就绪 →
sync.Once/close(ch) - 等待一组 →
WaitGroup - 高频计数 →
atomic
见 go-sync-package。谚语:不要通过共享内存通信;通过通信共享内存——但不是禁止 Mutex。
11. 背压(backpressure)
有界 channel 与有界 worker 把压力传回上游:
jobs := make(chan Job, 64) // 满则生产者阻塞
无界队列(无限 slice+mutex)会把内存当缓冲,延迟爆发。设计时明确:
- 阻塞
- 丢弃
- 拒绝(返回错误)
- 采样
12. Tee / 广播
一入多出:
- 对每个输出发送(注意取消)
- 或改用订阅模型
简单 tee 易在慢消费者处卡住整链——需策略。
13. 与内存模型
跨 goroutine 可见性依赖同步事件(unlock、channel 通信、WaitGroup 等)。模式里的 channel 操作通常已提供 happens-before。见 go-memory-model、go-happens-before-and-synchronization。
边界与限制
- 模式不能消除业务串行约束:有些步骤天生不能并行。
- 过度流水线:阶段间拷贝与调度开销可能比串行更慢。
- buffer 大小:经验值要压测,不是玄学常数。
- 取消后的残局:HTTP 客户端、DB 驱动是否真正中断,取决于实现。
- panic 在子 goroutine:默认整进程崩溃(Go 版本行为需注意);要隔离需自己 recover。
- 公平性:select 伪随机,不保证饥饿自由的业务公平。
- 分布式:这些是进程内模式;跨服务要队列/编排系统。
- 测试难度:时间与调度不确定性 → 用 ctx、假时钟、goleak、
-race。
常见误区
[!warning] 误区 1:开越多 goroutine 越快 受 CPU、IO、下游限额约束;无界 fan-out 常更慢更脆。
[!warning] 误区 2:到处 close channel 表示“完成” close 有所有权规则;完成也可用 WaitGroup+回调。见所有权笔记。
[!warning] 误区 3:超时后不管后台 goroutine 可能泄漏或继续写已丢弃的结果。
[!warning] 误区 4:用无缓冲 channel 做锁的替代却形成死锁 先画双人收发顺序。
[!warning] 误区 5:在循环里 time.After 做长时间服务 timer 压力与历史泄漏问题;优先 ctx 或复用 timer。
[!warning] 误区 6:共享 map 无锁只靠“我感觉并发不高” 用 race detector 打脸;见 go-race-detector。
[!warning] 误区 7:pipeline 每阶段无缓冲,却期望高吞吐 可能强串行握手;按需加 cap 或合并阶段。
[!warning] 误区 8:把 errgroup 当忽略错误的 WaitGroup 必须处理
Wait()错误。
工程实践
选择流程
是否真要并发? → 否:写清楚串行
需要取消? → context 从头穿到尾
任务数 vs 成本 → 有界并行 / 池
多阶段变换 → pipeline
多源合并 → fan-in + 单 closer
CPU 密集 → GOMAXPROCS 感知,少过度拆
IO 密集 → 有界 fan-out 常有效
代码审查清单
- 谁 close?
- 取消路径是否解除所有阻塞?
- 错误如何返回?丢失吗?
- 有无无界 goroutine 创建?
- 有无共享 map/slice 写?
- 测试是否
-race?
可观测性
- 队列长度、worker 忙碌、任务延迟
- ctx 取消原因计数
- 泄漏导致的 goroutine 数上涨(
runtime.NumGoroutine指标)
标准扩展库
官方博客路径
- Share Memory By Communicating
- Go Concurrency Patterns: Pipelines and cancellation
- Go Concurrency Patterns: Context
- Advanced Go Concurrency Patterns(历史 talk,概念仍有用)
实验
实验 A:pipeline 平方
运行最小示例,改 gen 不 close,观察 range 死锁。
实验 B:fan-in 唯一 closer
去掉 wg.Wait 前的独立 closer goroutine,改在某个 worker 里 close,制造 panic 或丢数据场景。
实验 C:有界 vs 无界
对 10000 个假 IO 任务分别:
- 每个任务一个 goroutine
- semaphore=8
比较最大 goroutine 数与总耗时(可用短 Sleep 模拟)。
实验 D:取消泄漏
package main
import (
"context"
"fmt"
"time"
)
func leaky(ctx context.Context) {
ch := make(chan int)
go func() { ch <- 1 }() // 可能无人收
select {
case <-ctx.Done():
return // 发送方可能泄漏
case v := <-ch:
fmt.Println(v)
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), time.Millisecond)
defer cancel()
leaky(ctx)
time.Sleep(20 * time.Millisecond)
fmt.Println("done")
}
修复:给发送方也监听 ctx.Done(),或缓冲 channel。
总结
| 问题 | 答案 |
|---|---|
| 模式解决什么? | 拓扑、完成信号、取消、有界与汇合 |
| 第一原则? | 每个 goroutine 有退出路径 |
| 取消用什么? | context 优先 |
| 关闭谁负责? | 发送所有权方 / 协调后的唯一 closer |
| 何时上 Mutex? | 共享内存状态比 channel 更直接时 |
| 与分布式? | 本篇是进程内;跨服务另议 |
一句话收束:
先保证能停、能关、能取消,再追求快;有界与所有权是并发模式的安全带。
自测题
1. pipeline 阶段为什么通常 defer close(out)?
2. fan-in 为何需要 WaitGroup + 单 closer?
3. 无界 fan-out 的主要风险?
4. 超时后仍泄漏的典型原因?
5. 有界 channel 如何形成背压?
6. errgroup.WithContext 相比普通 WaitGroup 多了什么?
答案
- 宣告该阶段不再发送,让下游 range 结束并连锁完成。
- 多发送者时必须等全部结束才能 close;且只能 close 一次。
- 资源耗尽、下游过载、调度与内存压力。
- 后台发送/接收仍阻塞在 channel 上,没有 done/ctx 分支。
- 缓冲满时发送者阻塞,从而减缓上游生产。
- 错误导致取消派生 ctx,兄弟任务可协作退出。