写并发代码最容易犯的错,是把「并发」当成「加速」。
大多数时候你需要的不是更快,而是让几件事互不阻塞。这两者的区别决定了你该用哪种模式——用错了,代码会变慢,而且更难调试。
下面三种模式覆盖了我在生产环境里遇到的绝大部分场景。

一、扇出扇入 #
要解决的问题: 手头有一批彼此独立的工作,想并行做完,再在一个地方汇总结果。
典型场景:批量请求多个下游接口、同时查缓存/DB/远程配置、爬虫并发抓 N 个 URL。
先看最干净的「多路汇入」 #
多个 channel 各自生产数据,合成一个输出 channel:
// fanIn 把多个只读 channel 的数据合并到一个输出 channel。
// 任一输入关掉后,对应 goroutine 退出;全部结束后再关闭 out。
func fanIn(inputs ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, in := range inputs {
wg.Add(1)
// 每个输入 channel 单独起一个「搬运工」
go func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v
}
}(in)
}
// 关键:必须另起一个 goroutine 等全部搬运结束再 close。
// 若在主流程里 Wait 再 close,主流程会卡死;若不 close,下游 range 永远不结束。
go func() {
wg.Wait()
close(out)
}()
return out
}
关键在最后那个 goroutine:没有它,out 永远不会关闭,下游的 range 会一直阻塞。
实战例子:并发请求多个 URL,汇总状态码 #
更贴近业务的是「扇出干活 + 扇入收结果」:主流程发任务、收结果,worker 池固定数量。
// URL 探测结果
type probeResult struct {
URL string
Status int
Err error
}
// probeURLs 并发探测 urls,最多 maxWorkers 个同时进行,返回与输入顺序无关的结果切片。
func probeURLs(ctx context.Context, urls []string, maxWorkers int) []probeResult {
if maxWorkers < 1 {
maxWorkers = 1
}
jobs := make(chan string)
results := make(chan probeResult)
// —— 扇出:固定数量的 worker 抢任务 ——
var wg sync.WaitGroup
for i := 0; i < maxWorkers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for u := range jobs {
// 每个任务自带 ctx,超时或取消时尽快退出
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
results <- probeResult{URL: u, Err: err}
continue
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
results <- probeResult{URL: u, Err: err}
continue
}
resp.Body.Close()
results <- probeResult{URL: u, Status: resp.StatusCode}
}
}()
}
// 任务发完后关闭 jobs,worker 才能退出
go func() {
for _, u := range urls {
select {
case <-ctx.Done():
close(jobs)
return
case jobs <- u:
}
}
close(jobs)
}()
// 全部 worker 结束后关闭 results,汇总方 range 才能结束
go func() {
wg.Wait()
close(results)
}()
// —— 扇入:在一个地方收齐结果 ——
var out []probeResult
for r := range results {
out = append(out, r)
}
return out
}
什么时候用: 任务之间没有先后依赖,只在乎「全部做完」或「尽量多做」。
什么时候别用: 下一步依赖上一步结果(那是流水线,不是简单扇出)。
二、用 context 取消 #
要解决的问题: 请求已经超时、用户关掉了页面、上游不需要结果了,后台 goroutine 还在跑——浪费资源,还可能在「已经没人听」的 channel 上发送导致泄漏。
任何可能长时间运行的 goroutine 都应该能被取消。
最小形态 #
select {
case <-ctx.Done():
// 超时、手动 cancel、父 context 取消,都走这里
return ctx.Err()
case result := <-work:
return handle(result)
}
实战例子:HTTP 接口里「调下游,最多等 2 秒」 #
// 查询用户订单:整体最多 2 秒;超时则返回错误,不再傻等下游。
func getOrderHandler(w http.ResponseWriter, r *http.Request) {
// 从请求派生超时 context;函数返回时 cancel,通知所有子调用停手
ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second)
defer cancel()
orderID := r.URL.Query().Get("id")
order, err := fetchOrder(ctx, orderID)
if err != nil {
// ctx 超时常见错误:context.DeadlineExceeded
http.Error(w, err.Error(), http.StatusGatewayTimeout)
return
}
_ = json.NewEncoder(w).Encode(order)
}
// fetchOrder 模拟访问下游;必须把 ctx 传到真正会阻塞的地方(HTTP、DB、RPC)。
func fetchOrder(ctx context.Context, id string) (map[string]string, error) {
// 模拟异步结果
ch := make(chan map[string]string, 1)
go func() {
// 假装下游很慢
time.Sleep(5 * time.Second)
// 若已取消,最好别再往外送(这里用 buffer channel 避免泄漏演示变复杂)
select {
case <-ctx.Done():
return
case ch <- map[string]string{"id": id, "status": "paid"}:
}
}()
select {
case <-ctx.Done():
// 2 秒到了:调用方立刻返回,不必等满 5 秒
return nil, ctx.Err()
case order := <-ch:
return order, nil
}
}
再补一个日常会写的「可取消的循环任务」:
// 后台同步:每隔 interval 拉一次数据,直到 ctx 取消(进程退出 / 配置热更停掉任务)。
func syncLoop(ctx context.Context, interval time.Duration) {
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
log.Println("同步任务已停止:", ctx.Err())
return
case <-t.C:
if err := pullOnce(ctx); err != nil {
log.Println("本轮同步失败:", err)
}
}
}
}
习惯:
- 入口造
ctx(HTTP 用r.Context()+WithTimeout); - 往下传,不要另起无关联的
context.Background(); - 在
select/ 库函数参数里真正响应ctx.Done()。
三、信号量限流 #
要解决的问题: 无限制地开 goroutine 是新手最常见的坑——瞬时几千个请求下游、打满文件描述符、拖垮数据库连接池。
用带缓冲的 channel 当信号量:缓冲区大小 = 最大并发数。
最小形态 #
// 最多 10 个任务同时执行
sem := make(chan struct{}, 10)
for _, task := range tasks {
sem <- struct{}{} // 池子满了就阻塞在这里,不会无限开 goroutine
go func(t Task) {
defer func() { <-sem }() // 做完归还名额
process(t)
}(task)
}
注意:上面循环结束后,主 goroutine 不会等全部 process 结束。若要等齐,再加 WaitGroup。
实战例子:批量写库,最多 5 个并发 #
// saveUsers 把一批用户写入存储;同一时刻最多 5 个写操作在飞。
func saveUsers(ctx context.Context, users []User) error {
const maxConcurrent = 5
sem := make(chan struct{}, maxConcurrent)
var wg sync.WaitGroup
// 用 errgroup 或下面这个「第一个错误记下、其余继续」的简化版
errCh := make(chan error, 1)
for _, u := range users {
// 先响应取消,避免还在往池子里塞任务
select {
case <-ctx.Done():
wg.Wait()
return ctx.Err()
case sem <- struct{}{}:
}
wg.Add(1)
go func(user User) {
defer wg.Done()
defer func() { <-sem }() // 释放并发名额
if err := dbInsert(ctx, user); err != nil {
// 只保留第一个错误即可
select {
case errCh <- err:
default:
}
}
}(u)
}
wg.Wait()
select {
case err := <-errCh:
return err
default:
return nil
}
}
和扇出扇入怎么配合? #
限流管的是「同时跑几个」;扇出扇入管的是「怎么分发、怎么汇总」。
实战里常常叠在一起:信号量限制 worker 数量,channel 传递任务和结果——第一节 probeURLs 里的 maxWorkers 本质上就是限流。
经验值:
- 调外部 HTTP:并发常从 5~20 试起;
- 写本地磁盘:要看磁盘与锁;
- 打数据库:不要超过连接池大小。
小结 #
| 模式 | 一句话 | 常见用途 |
|---|---|---|
| 扇出扇入 | 分头干活,再汇总 | 批量请求、多源查询、并发探测 |
| context 取消 | 随时喊停 | 接口超时、后台任务停机、级联取消 |
| 信号量限流 | 同时最多 N 个 | 保护 DB/下游、批量导入、爬虫 |
三种模式解决的是三个不同的问题,不要混用:
需要汇总用扇出扇入;需要停用 context;需要别把系统打挂用限流。
生产里多数接口是「context + 有限 worker 的扇出」,而不是裸 go func 一把梭。