Go 并发模式概览

常见进程内并发结构:pipeline、fan-out/fan-in、超时与取消、or-done、防止泄漏;与 worker pool 专题的分工。

#type / concept #status / growing #tech / dev #resource / go

[!info] 关联笔记

Go 并发模式概览

这个概念为什么会出现

掌握 goroutine、channel、select 之后,工程里反复撞上同一类结构问题

  • 多阶段处理如何串成流水线?
  • 如何扇出并行又限制并发度?
  • 多个结果如何汇合?
  • 超时/取消时如何不泄漏 goroutine?
  • 谁关闭 channel?完成信号如何传播?

这些问题不是再背一个 API 能解决的,而是可复用的进程内并发拓扑。本篇给模式地图与判定标准;具体落地可下钻 go-worker-poolgo-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

结合场景再看三个关注点

  1. 发送方关闭,接收方 range
    订单流结束靠 close,不是另传“EOF 业务消息”(除非协议需要)。

  2. 阶段函数签名暴露只读 in/out
    所有权清晰,避免“谁都可以 close”的竞态。

  3. 超时是并发模式的一部分
    支付回调用 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)

原则:

  1. 阶段内拥有 out 的发送与关闭权
  2. 上游关闭 → 本阶段 range 结束 → 关闭下游
  3. 数据与完成信号都走 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
}

关键:

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
}

优先级建议:

  1. context 贯穿全链(请求级取消)
  2. 循环内慎用 time.After(每次分配 timer,历史版本易泄漏等待中的 timer;可用 time.NewTimerStop
  3. 超时后仍要考虑后台 goroutine 是否退出

go-contextgo-select

6. 防止 goroutine 泄漏

泄漏典型形态:

  • 发送到无人接收且无缓冲 channel
  • range 未关闭的 channel
  • 等待 ctx 但从不知道 cancel
  • select 缺少 done 分支

检查清单:

  1. 每个 go 的退出条件是什么?
  2. 取消时阻塞在 channel 上的双方能否解锁?
  3. 测试里用 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
  • 收集所有错误 → 自己设计 []errorerrors.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-modelgo-happens-before-and-synchronization

边界与限制

  1. 模式不能消除业务串行约束:有些步骤天生不能并行。
  2. 过度流水线:阶段间拷贝与调度开销可能比串行更慢。
  3. buffer 大小:经验值要压测,不是玄学常数。
  4. 取消后的残局:HTTP 客户端、DB 驱动是否真正中断,取决于实现。
  5. panic 在子 goroutine:默认整进程崩溃(Go 版本行为需注意);要隔离需自己 recover。
  6. 公平性:select 伪随机,不保证饥饿自由的业务公平。
  7. 分布式:这些是进程内模式;跨服务要队列/编排系统。
  8. 测试难度:时间与调度不确定性 → 用 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 常有效

代码审查清单

  1. 谁 close?
  2. 取消路径是否解除所有阻塞?
  3. 错误如何返回?丢失吗?
  4. 有无无界 goroutine 创建?
  5. 有无共享 map/slice 写?
  6. 测试是否 -race

可观测性

  • 队列长度、worker 忙碌、任务延迟
  • ctx 取消原因计数
  • 泄漏导致的 goroutine 数上涨(runtime.NumGoroutine 指标)

标准扩展库

官方博客路径

  1. Share Memory By Communicating
  2. Go Concurrency Patterns: Pipelines and cancellation
  3. Go Concurrency Patterns: Context
  4. Advanced Go Concurrency Patterns(历史 talk,概念仍有用)

实验

实验 A:pipeline 平方

运行最小示例,改 genclose,观察 range 死锁。

实验 B:fan-in 唯一 closer

去掉 wg.Wait 前的独立 closer goroutine,改在某个 worker 里 close,制造 panic 或丢数据场景。

实验 C:有界 vs 无界

对 10000 个假 IO 任务分别:

  1. 每个任务一个 goroutine
  2. 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 多了什么?

答案
  1. 宣告该阶段不再发送,让下游 range 结束并连锁完成。
  2. 多发送者时必须等全部结束才能 close;且只能 close 一次。
  3. 资源耗尽、下游过载、调度与内存压力。
  4. 后台发送/接收仍阻塞在 channel 上,没有 done/ctx 分支。
  5. 缓冲满时发送者阻塞,从而减缓上游生产。
  6. 错误导致取消派生 ctx,兄弟任务可协作退出。

依据与延伸

创建于 2026/7/14 更新于 2026/7/15