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. 错误传播的三种策略

并发组通常选择以下语义之一:

  1. 首错取消:任何关键任务失败就取消同组任务,等待全部退出,返回首个或带原因的错误。
  2. 收集全部:任务彼此独立,完成所有任务后返回结果和错误集合。
  3. 容忍部分失败:记录失败并继续,但必须定义成功阈值、重试和最终状态。

不能只让 worker 写一个无缓冲 errCh 然后立即退出:多个错误可能让后续 worker 卡在发送,拥有者也可能在收到首错后不再消费。首错模式可用容量 1 的错误 channel 配合非阻塞发送和 cancel;仍需 Wait 确认全部 worker 退出。收集全部模式可由单一聚合器消费结构化结果,避免多个 goroutine append 同一个 slice 的竞态。

标准库 WaitGroup 不处理 error。项目若使用 golang.org/x/sync/errgroupWithContextSetLimitTryGo 能表达首错取消及并发限制;但它仍不替你定义队列、部分结果、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 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。