跳到主要内容

Go 高并发实战(2):结构化并发、Context 与有界 Worker Pool

Rainy
雨落无声,代码成诗 —— 致力于技术与艺术的极致平衡
Rainy
8 MIN READ... VIEWS

创建 goroutine 的代码,也必须能说明谁等待它、谁取消它、它什么时候退出。

本期把 goroutine 从语法糖升级为受管理的生命周期。我们实现两种核心结构:请求内 fan-out/fan-in,以及跨请求的有界 worker pool。

架构图:父请求必须拥有全部子任务

图里最重要的不是三条并行分支,而是闭合的生命周期:父作用域创建子任务、传播取消、收集错误并等待退出。任何绕过 Join 的 goroutine 都可能在 handler 返回后继续持有连接、锁或引用对象。

一、结构化并发的四个问题

每次写 go fn() 前回答:

  1. Owner 是谁:哪个作用域负责它?
  2. Join 在哪:谁等待完成并观察错误?
  3. Cancel 从哪来:调用方离开后如何停止?
  4. 资源上限是多少:最坏会同时存在多少个 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。建议:

  1. 预热服务并记录稳定基线;
  2. 重复运行超时/取消路径数百次;
  3. 等待一个明确 grace period;
  4. 比较 goroutine profile 中业务栈,而不仅是总数;
  5. 失败时保存完整 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 与数据库并发

参考资料

Logo
RainLib

探索技术、设计与分布式系统的边界。构建面向未来的开发者工具。

留言与建议

© 2026 RainLib. 为未来构建。(Built for the Future)
版权所有。
系统正常