Go 基础体系 · 第 21/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go 并发模式:Worker Pool、Pipeline、Fan-out 与背压
本文以 Go 1.26.4 为基准。并发模式不是固定代码模板,而是对资源和生命周期的约束:谁接纳工作、同时最多执行多少、排队能占多少内存、失败是否取消同组任务、结果如何汇聚、关停时谁等待。Worker Pool、Pipeline、Fan-out/Fan-in 都建立在这些约束之上;少了取消与关闭协议,示例看似能跑,生产中却会泄漏或在过载时失控。
本篇专注端到端拓扑。channel 的关闭与通信细节、select/Timer、锁和内存模型各有相邻主题,但这里会完整说明模式成立所需的协议,不用链接代替关键推理。
1. 先判断是否真的需要并发
并发能重叠独立等待或利用多个 CPU 核,却会增加调度、同步、内存和故障状态。对很小的 CPU 任务,创建 goroutine 与传递结果可能比工作本身更贵;对数据库调用,并发度超过连接池或数据库容量只会扩大排队。
设计前量化四个量:到达率、单项服务时间、可接受等待、下游并发容量。Little's Law 给出稳定系统的直觉:系统内平均任务数约等于到达率乘平均停留时间。若到达率长期高于处理能力,任何有限队列最终都会满,任何无限队列最终都会耗尽内存。模式不能创造容量,只能让过载决策显式发生。
顺序循环通常是正确性基线。先写出可测试的单项函数,再决定在哪一层并行;这样错误重试、指标和资源释放不会散落到每个 goroutine。
2. Worker Pool 的完整生命周期
固定 worker pool 用一组长期 goroutine 消费共享任务 channel。典型角色是:提交者发送并最终关闭 jobs;worker 读取到关闭,处理并发送结果;协调者等待所有 worker 后关闭 results;拥有者消费结果并处理取消。
producer --jobs--> worker 1 --+
worker 2 ---+--results--> owner
worker N ---+
wait then close results
worker 数限制同时执行量,jobs 容量限制已接纳但未执行的量,两者不是同一个参数。输出 channel 的关闭者必须是知道所有 worker 已退出的协调者,不能让任意 worker 在自己结束时关闭共享输出。
Go 1.26.4 可用 WaitGroup.Go 管理正常返回的 worker;传入函数不应 panic。若业务处理可能 panic,应在明确的任务边界恢复、转换为错误并确保系统状态仍可信,不能把 recover 当普通错误处理。
3. 有界队列与四种背压决策
队列满时只有几类选择,每种都应由业务而非偶然调度决定:
- 阻塞:把压力传回调用方,适合必须处理且上游允许等待;阻塞必须可取消。
- 拒绝:立即返回 overload,适合在线请求;调用方可按预算重试或降级。
- 丢弃:适合允许损失的遥测或刷新信号;要明确丢最新、丢最旧还是合并。
- 溢出到持久层:适合必须最终处理的任务,但这已经是持久队列协议,需要幂等、确认和重放。
func submit(ctx context.Context, jobs chan<- Job, job Job) error {
select {
case jobs <- job:
return nil
case <-ctx.Done():
return context.Cause(ctx)
}
}
用无界 slice 加一个通知 channel 只是把背压变成内存风险。队列容量应按最大内存和允许等待反推,并监控队列等待时间而不只是长度。同一长度在处理速度变化时代表完全不同的用户延迟。
4. 并发限制不等于速率限制
worker 数或令牌信号量限制同一时刻在途任务数,不能限制每秒请求数。若每项从 10 ms 变成 1 s,固定并发会自动降低吞吐,这是保护资源的效果;但任务很快时它可能瞬间发出大量请求。
速率限制控制时间窗口内启动数量,通常使用成熟令牌桶实现。生产调用下游时常常两者都需要:并发上限保护连接、内存和 goroutine;速率上限满足对方配额。把 ticker 每次发一个 token 的简易实现用于严格配额时要谨慎处理突发量、时钟漂移和取消。
数据库场景还应与 SetMaxOpenConns 协调。若应用有 100 个 worker 而连接池只有 10,额外 90 个只是在另一处排队;更糟的是任务可能持有其他资源再等待连接,扩大耦合。
5. Pipeline:每个 stage 都要能结束
Pipeline 把处理分为多个 stage,每个 stage 从输入读取、转换并写入自己拥有的输出。一个稳健 stage 有三条规则:只关闭自己的输出;输入关闭后处理完再返回;所有可能阻塞的输出都能观察取消。
func mapStage[A, B any](
ctx context.Context,
input <-chan A,
fn func(A) B,
) <-chan B {
output := make(chan B)
go func() {
defer close(output)
for value := range input {
mapped := fn(value)
select {
case output <- mapped:
case <-ctx.Done():
return
}
}
}()
return output
}
如果下游只读取第一个结果便返回,上游 stage 可能永久阻塞在发送。因此取消必须贯穿整条 pipeline,而且最终拥有者退出前要调用 cancel。只在 stage 的循环顶部检查 context 不够,阻塞发送本身也要放进 select。
每个 stage 一个 goroutine 不是硬规则。CPU 密集 stage 可内部 fan-out,纯粹的轻量转换可以合并,避免多次调度和复制。stage 边界应该对应独立的并发度、缓冲、错误或观测需求。
6. Fan-out:并行消费与工作分配
多个 worker 从同一个 input channel 接收是最简单的 fan-out。每个元素只交给其中一个 worker,适合任务相互独立且允许乱序。它不是广播;所有订阅者都必须收到时,需要为每个订阅者维护投递状态与慢消费者策略。
并行度选择取决于任务类型:CPU 密集工作通常从 GOMAXPROCS 附近试验;I/O 工作可更高,但受连接池、文件描述符、对方限流和每项内存约束。动态调大 worker 不能修复慢下游,只可能使其更过载。
任务函数应接收 context,并把它传给 HTTP、数据库等可取消 API。若调用的是不可中断的第三方函数,worker 收到取消也只能等它返回;为每项再开一个 goroutine 会把等待转移成泄漏。要从 API 或隔离边界解决不可取消问题。
7. Fan-in:唯一协调者关闭输出
Fan-in 合并多个输入到一个输出。每个转发 goroutine 结束后递减 WaitGroup,独立协调者等待并关闭输出:
func merge[T any](ctx context.Context, inputs ...<-chan T) <-chan T {
output := make(chan T)
var wg sync.WaitGroup
for _, input := range inputs {
input := input
wg.Go(func() {
for value := range input {
select {
case output <- value:
case <-ctx.Done():
return
}
}
})
}
go func() {
wg.Wait()
close(output)
}()
return output
}
如果某个 input 从不关闭且 context 不取消,合并输出也不会关闭。这正是生命周期协议的一部分,不是 merge 能自行猜测的异常。输入数量极大时,每输入一个 goroutine 会有成本,可由单调度循环或分层合并降低,但复杂度应由 profile 和规模证明。
合并结果的相对顺序不确定;同一输入内部的发送顺序通常保持,但来自不同输入的交错不应依赖。
8. 顺序、重排与队头阻塞
并行 worker 天然可能乱序完成。若 API 要保持输入顺序,可以给任务携带递增序号,聚合器用 map 暂存超前结果,等待下一个序号再连续输出。
这会引入队头阻塞:序号 10 很慢时,11 到 100 即使完成也必须占内存等待。需要规定最大重排窗口、单项 deadline 和失败时是否跳过。另一种方法是把输入分片到多个有序 lane,例如同一用户键固定进入同一 worker;它保留键内顺序而允许键间并行,但热点键可能形成倾斜。
不要不加说明地“按完成顺序返回”,也不要为了全局顺序无限缓存。顺序是 API 契约,会直接影响吞吐、内存和尾延迟。
9. 错误传播的三种策略
并发组通常选择以下语义之一:
- 首错取消:任何关键任务失败就取消同组任务,等待全部退出,返回首个或带原因的错误。
- 收集全部:任务彼此独立,完成所有任务后返回结果和错误集合。
- 容忍部分失败:记录失败并继续,但必须定义成功阈值、重试和最终状态。
不能只让 worker 写一个无缓冲 errCh 然后立即退出:多个错误可能让后续 worker 卡在发送,拥有者也可能在收到首错后不再消费。首错模式可用容量 1 的错误 channel 配合非阻塞发送和 cancel;仍需 Wait 确认全部 worker 退出。收集全部模式可由单一聚合器消费结构化结果,避免多个 goroutine append 同一个 slice 的竞态。
标准库 WaitGroup 不处理 error。项目若使用 golang.org/x/sync/errgroup,WithContext、SetLimit 和 TryGo 能表达首错取消及并发限制;但它仍不替你定义队列、部分结果、panic 和外部 I/O 取消。
10. 取消不是完成,关停需要等待
调用 cancel() 只是关闭一个信号,不表示 goroutine 已退出。优雅关停应按依赖方向执行:停止接纳新任务;取消或关闭生产入口;等待已接纳任务完成或达到关停 deadline;关闭由本层拥有的资源;最后确认 worker 全部结束。
stop admission -> signal cancellation -> workers stop/finish
-> wait group returns -> close result stream -> release resources
关停预算到期后需要明确策略。进程退出会直接终止 goroutine,defer 不保证执行;因此必须持久化的任务不能只存在内存队列。对于 HTTP 服务,先让服务器停止新连接并等待 handler,再停止其共享后台组件,避免 handler 使用已关闭资源。
每个启动 goroutine 的位置都应能回答:谁取消、谁等待、错误去哪。库函数若返回 channel 并偷偷启动 goroutine,应在 API 中说明调用者停止消费时如何让它退出。
11. 批处理、分片与工作窃取
批处理把多项合并成一次数据库、RPC 或 channel 操作,可提高吞吐,却增加第一项的等待时间和失败重试范围。常用触发条件是“达到最大条数或最长等待时间”,两者都必须有界。取消时是刷新尾批还是丢弃,要按持久性契约决定。
按键分片可减少共享锁并保持键内顺序。分片函数必须稳定,扩缩分片会改变映射并可能破坏顺序;一致性哈希或迁移协议属于更高层问题。工作窃取能均衡不均匀任务,但实现复杂,Go 调度器已经会在 goroutine 层做调度;应用通常先使用共享 jobs channel。
为每个请求建立一个 worker pool 通常浪费。服务级共享池能统一限流,但要防止一个租户占满队列;可使用每租户配额、公平队列或分层信号量。
12. 常见失败模式
生产中反复出现的错误包括:
- 无界创建 goroutine,把突发流量转换成栈、连接和调度压力。
- 无界队列掩盖过载,直到 OOM 或任务过期才失败。
- 下游提前返回,上游仍阻塞发送,goroutine 数持续增长。
- worker 关闭共享结果 channel,与其他发送者竞争并 panic。
- 收到首错就返回,没有 cancel 和 Wait,其余任务仍运行。
- 在多个 goroutine 中直接 append 结果 slice 或写普通 map,造成竞态。
- 令牌在错误路径未归还,实际并发度逐步降为零。
- 把并发上限误作每秒速率,触发下游限流。
- 重试任务重新入队但没有次数与总 deadline,形成自激流量。
缓冲加大、Sleep 或 recover 只能改变故障出现时机。修复应回到所有权、容量和状态转换。
13. 诊断并发拓扑
最有用的指标能对应队列生命周期:接纳数、拒绝数、当前/最大队列长度、排队时间、活跃 worker、执行时间、成功/失败/取消数、重试数和关停耗时。吞吐平稳但排队时间增长说明系统正在积债;CPU 不高也可能是下游或锁容量不足。
go test -race ./...
go test -run TestPipeline -count=100 ./...
go tool pprof http://127.0.0.1:6060/debug/pprof/goroutine
go tool pprof http://127.0.0.1:6060/debug/pprof/block
go tool trace trace.out
goroutine profile 中大量相同 [chan send] 通常指向下游停止或队列已满;[select] 要结合 context 是否还能取消。block profile 展示累计阻塞热点,trace 可还原调度与任务阶段。抓取多份 profile 看增长趋势,不把正常常驻 worker 误判为泄漏。
日志应携带任务 ID、尝试次数、队列等待和取消原因,但避免每项高频阶段都打印。分布式 trace 可把排队 span 与执行 span 分开,否则“下游慢”可能实际是本地等队列。
14. 确定性测试与故障注入
并发测试用 channel barrier 控制事件,不用 Sleep 猜顺序。至少覆盖:空输入、单项、任务数少于/多于 worker、处理失败、调用方取消、输出消费者提前停止、队列满、关停 deadline 和 panic 策略。
started := make(chan struct{})
release := make(chan struct{})
process := func(ctx context.Context, job Job) error {
close(started)
select {
case <-release:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
测试看到 started 后再触发取消,能精确覆盖执行中取消。测试最终要有宽松 watchdog,失败时输出 goroutine dump。使用 -race 和重复次数扩大交错覆盖,但不要通过把超时不断调大隐藏泄漏。
基准测试要区分单项函数成本、调度开销和端到端吞吐,使用真实任务分布而非全部等时任务。分别记录 ops/s、P95/P99、allocs/op 和最大内存;平均吞吐提升可能以更差队头阻塞为代价。
15. 生产调参顺序
先确定下游安全并发与内存预算,再设 worker 和队列;压测稳定负载与突发;观察排队和执行分位数;最后才改变批量、分片或动态伸缩。每次只改变一个主要约束,并保留拒绝或降级路径。
动态 worker 数需要防振荡:扩容依据队列延迟而非瞬时长度,缩容要有冷却期,且永远不突破下游硬容量。很多服务用固定池加自动水平扩容已经足够,进程内再自适应会与外部扩容器相互干扰。
容量规划还要包含每个排队任务引用的数据大小、worker 栈、结果缓冲、重试副本和外部连接。只计算 channel 槽位会严重低估内存。
16. 可运行综合示例:首错取消的有序 Worker Pool
下面程序实现有界 jobs、固定 worker、首错取消、完整等待以及按输入序号重排。拥有者始终消费到 results 关闭,因此 worker 不会卡在结果发送。首个错误触发取消,但已经成功提交的结果仍可能到达;最终返回错误并丢弃部分成功,由 API 明确定义。
package main
import (
"context"
"errors"
"fmt"
"sync"
)
type job struct {
index int
value int
}
type result struct {
index int
value int
err error
}
func orderedSquares(ctx context.Context, values []int, workers int) ([]int, error) {
if workers < 1 {
return nil, errors.New("workers must be positive")
}
ctx, cancel := context.WithCancelCause(ctx)
defer cancel(nil)
jobs := make(chan job, workers)
results := make(chan result, workers)
go func() {
defer close(jobs)
for index, value := range values {
select {
case jobs <- job{index: index, value: value}:
case <-ctx.Done():
return
}
}
}()
var wg sync.WaitGroup
for range workers {
wg.Go(func() {
for item := range jobs {
current := result{index: item.index}
if item.value < 0 {
current.err = fmt.Errorf("negative value at %d", item.index)
} else {
current.value = item.value * item.value
}
select {
case results <- current:
if current.err != nil {
cancel(current.err)
return
}
case <-ctx.Done():
return
}
}
})
}
go func() {
wg.Wait()
close(results)
}()
output := make([]int, len(values))
completed := 0
for current := range results {
if current.err == nil {
output[current.index] = current.value
completed++
}
}
if cause := context.Cause(ctx); cause != nil {
return nil, cause
}
if completed != len(values) {
return nil, fmt.Errorf("completed %d of %d jobs", completed, len(values))
}
return output, nil
}
func main() {
values, err := orderedSquares(context.Background(), []int{5, 2, 9, 3}, 3)
if err != nil {
panic(err)
}
fmt.Println(values)
}
这个实现追求协议清晰,不声称是所有负载的最快方案。生产版本还应给外层 context 设置 deadline,记录排队与执行耗时,并根据“部分结果是否有价值”决定错误返回结构。并发模式的成熟标志不是 goroutine 数量,而是正常、失败、取消和过载四条路径都能有界结束。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go sync 与 atomic:Mutex、RWMutex、WaitGroup、Once 和 Cond
- 下一篇:Go context 完整指南:取消、超时、Deadline 与 Value
- 延伸:Go channel 完整基础:发送、接收、缓冲、关闭与所有权
- 延伸:Go goroutine 生命周期:泄漏、打断、错误传播与优雅关闭
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论