Go 高并发实战(2):结构化并发、Context 与有界 Worker Pool
创建 goroutine 的代码,也必须能说明谁等待它、谁取消它、它什么时候退出。
本期把 goroutine 从语法糖升级为受管理的生命周期。我们实现两种核心结构:请求内 fan-out/fan-in,以及跨请求的有界 worker pool。
架构图:父请求必须拥有全部子任务
图里最重要的不是三条并行分支,而是闭合的生命周期:父作用域创建子任务、传播取消、收集错误并等待退出。任何绕过 Join 的 goroutine 都可能在 handler 返回后继续持有连接、锁或引用对象。
一、结构化并发的四个问题
每次写 go fn() 前回答:
- Owner 是谁:哪个作用域负责它?
- Join 在哪:谁等待完成并观察错误?
- Cancel 从哪来:调用方离开后如何停止?
- 资源上限是多少:最坏会同时存在多少个 goroutine、job 和结果?
无法回答的 detached goroutine 是线上泄漏的高风险来源。
二、Context 是取消树,不是参数袋
入口 r.Context() 会在客户端断开或 handler 返回时取消。子调用继续派生更短 deadline:
func fetch(ctx context.Context, client *http.Client, url string) ([]byte, error) {
child, cancel := context.WithTimeout(ctx, 80*time.Millisecond)
defer cancel()
req, err := http.NewRequestWithContext(child, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
response, err := client.Do(req)
if err != nil {
return nil, err
}
defer response.Body.Close()
return io.ReadAll(io.LimitReader(response.Body, 1<<20))
}
规则:Context 作为第一个参数向下传递;不要存进长期对象;不要用 context.Background() 切断请求取消;每个 WithTimeout/WithCancel 都调用 cancel 释放 timer 资源。
三、请求内 fan-out/fan-in
type Result struct {
Source string
Value string
}
func Aggregate(ctx context.Context, sources []string) ([]Result, error) {
ctx, cancel := context.WithCancel(ctx)
defer cancel()
resultCh := make(chan Result, len(sources))
errCh := make(chan error, len(sources))
var wg sync.WaitGroup
wg.Add(len(sources))
for _, source := range sources {
source := source
go func() {
defer wg.Done()
value, err := fetchOne(ctx, source)
if err != nil {
errCh <- fmt.Errorf("%s: %w", source, err)
return
}
resultCh <- Result{Source: source, Value: value}
}()
}
done := make(chan struct{})
go func() {
wg.Wait()
close(done)
}()
results := make([]Result, 0, len(sources))
for len(results) < len(sources) {
select {
case result := <-resultCh:
results = append(results, result)
case err := <-errCh:
cancel()
return nil, err
case <-done:
if len(results) != len(sources) {
return nil, errors.New("incomplete aggregation")
}
case <-ctx.Done():
return nil, ctx.Err()
}
}
return results, nil
}
结果/错误 channel 容量覆盖所有生产者,因此聚合方提前返回后,生产者不会永久卡在发送。真正下游仍必须接受 context;buffer 只能防发送泄漏,不能中断不可取消 I/O。
四、跨请求 Worker Pool
type Job struct {
Ctx context.Context
Payload []byte
Reply chan<- error
}
type Pool struct {
jobs chan Job
wg sync.WaitGroup
}
func NewPool(workers, queue int) *Pool {
pool := &Pool{jobs: make(chan Job, queue)}
pool.wg.Add(workers)
for workerID := 0; workerID < workers; workerID++ {
go func() {
defer pool.wg.Done()
for job := range pool.jobs {
select {
case <-job.Ctx.Done():
job.Reply <- job.Ctx.Err()
default:
job.Reply <- process(job.Ctx, job.Payload)
}
}
}()
}
return pool
}
func (p *Pool) Submit(ctx context.Context, payload []byte) error {
reply := make(chan error, 1)
job := Job{Ctx: ctx, Payload: payload, Reply: reply}
select {
case p.jobs <- job:
case <-ctx.Done():
return ctx.Err()
}
select {
case err := <-reply:
return err
case <-ctx.Done():
return ctx.Err()
}
}
func (p *Pool) Close() {
close(p.jobs)
p.wg.Wait()
}
4.1 Pool 大小怎么定
- CPU 任务:worker 从 CPU 核数附近开始;
- 阻塞 I/O:由下游连接/并发容量决定;
- queue:
可处理速率 × 最大排队时间,不是“内存能放多少”; - 满队列:等待、拒绝、降级或持久化必须显式选择。
多个生产者仍在提交时不能直接 Close,否则会 panic。生产入口应先停,再由唯一 owner 关闭 channel,最后等待 worker。
五、典型泄漏模式
// 错误:调用方超时后没有接收者,发送永久阻塞。
result := make(chan Value)
go func() { result <- slowCall() }()
select {
case value := <-result:
return value
case <-ctx.Done():
return Value{}
}
检查清单:
- 每个阻塞 send/receive 是否同时监听取消;
- 结果 channel 是否能容纳提前返回后的发送;
time.NewTicker是否 Stop;- 下游 body/rows 是否 Close;
- goroutine 是否等待永不关闭的 channel;
- fan-out 数量是否受输入硬上限控制。
六、Channel 的所有权与关闭协议
channel close 是“不会再有新值”的广播事实,不是通用取消按钮。最稳妥的约定是:生产者关闭自己唯一拥有的输出 channel,接收者不关闭输入 channel。
func transform(ctx context.Context, in <-chan Job) <-chan Result {
out := make(chan Result)
go func() {
defer close(out)
for {
select {
case <-ctx.Done():
return
case job, ok := <-in:
if !ok {
return
}
result := run(job)
select {
case out <- result:
case <-ctx.Done():
return
}
}
}
}()
return out
}
value, ok := <-ch 区分零值与关闭;向已关闭 channel 发送会 panic,重复 close 也会 panic。多个生产者共享输出时,由一个协调 goroutine 在所有 producer Wait 后关闭,而不是让任意 producer 抢着 close。
6.1 Buffer 的三种意义
- 0:发送和接收 rendezvous,天然同步;
- 1:允许一次状态通知不依赖接收时机;
- N:允许生产/消费短暂错峰,同时形成 N 个元素的内存与延迟上限。
buffer 不改变长期处理能力。当平均生产速率大于消费速率,任何有限 buffer 最终都会满。
七、select、Timer 与取消边界
多个 case 同时 ready 时,不能依赖固定业务优先级。若 shutdown 必须优先,先做一次非阻塞取消检查,再进入主 select,或者将状态机写得对任意选择都正确。
循环中不要反复 time.After 创建新 timer:
timer := time.NewTimer(wait)
defer timer.Stop()
select {
case <-timer.C:
return ErrTimeout
case result := <-resultCh:
return result
case <-ctx.Done():
return ctx.Err()
}
复用 timer 时要正确 Stop、必要时排空 channel,再 Reset。timer bug 往往只在高频重试或长时间运行后表现为分配增长和幽灵超时。
八、WaitGroup 与错误传播
WaitGroup 只表达完成数量,不携带错误或取消。调用 Add 必须发生在启动 goroutine 之前,避免 goroutine 先 Done;不要在一次 Wait 周期未结束时危险复用。
复杂请求可以使用 errgroup.WithContext:第一个错误取消同组任务,Wait 汇总生命周期。但仍要:
- 设置 fan-out 上限;
- 确认任务 API 遵守 Context;
- 决定第一个错误是否真的应取消全部;
- 保留具体 source 与错误分类。
九、Goroutine 泄漏测试
泄漏测试不能只比较瞬间数量,因为 runtime 和测试框架也会创建后台 goroutine。建议:
- 预热服务并记录稳定基线;
- 重复运行超时/取消路径数百次;
- 等待一个明确 grace period;
- 比较 goroutine profile 中业务栈,而不仅是总数;
- 失败时保存完整
debug=2栈作为测试产物。
故障注入表:下游永不返回、结果消费者提前退出、队列已满、worker panic、关闭与 Submit 并发。每条路径都要有 owner、取消和 join 证据。
9.1 逐步复现“超时后仍工作”
func TestAggregateCancellation(t *testing.T) {
started := make(chan struct{})
stopped := make(chan struct{})
downstream := func(ctx context.Context) error {
close(started)
defer close(stopped)
<-ctx.Done()
return ctx.Err()
}
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond)
defer cancel()
err := runAggregate(ctx, downstream)
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("want deadline, got %v", err)
}
<-started
select {
case <-stopped:
case <-time.After(100 * time.Millisecond):
t.Fatal("child goroutine did not observe cancellation")
}
}
测试直接验证子任务接收取消并走到退出点,而不是“睡一秒再数 goroutine”。实际服务还应比较 goroutine profile,覆盖未被该 hook 捕获的泄漏路径。
9.2 Worker Pool 的背压时间线
队列容量控制“等待中的任务”,worker 数控制“正在执行的任务”,下游 pool 控制“真实稀缺资源”。三者必须分别监控,不能只暴露一个 concurrency 数字。
十、本期实战验收
- 对 worker 数、queue 长度和提交 deadline 建参数矩阵;
- 压到 queue 满,确认拒绝可观测且 goroutine/RSS 不持续增长;
- 在每个阶段随机取消请求,确认子任务与许可回收;
- 停止输入并调用
Close,用 goroutine profile 确认 worker 全部退出。
上一篇:Go-1:GMP 调度器与容量模型 · 下一篇:Go-3:生产级 HTTP、RPC 与数据库并发。