清风徐来
知是行之始,行是知之成

Go 的三种并发模式

Go Go并发 约 7 分钟
本文目录

写并发代码最容易犯的错,是把「并发」当成「加速」。

大多数时候你需要的不是更快,而是让几件事互不阻塞。这两者的区别决定了你该用哪种模式——用错了,代码会变慢,而且更难调试。

下面三种模式覆盖了我在生产环境里遇到的绝大部分场景。

一、扇出扇入 #

要解决的问题: 手头有一批彼此独立的工作,想并行做完,再在一个地方汇总结果。

典型场景:批量请求多个下游接口、同时查缓存/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)
			}
		}
	}
}

习惯:

  1. 入口造 ctx(HTTP 用 r.Context() + WithTimeout);
  2. 往下传,不要另起无关联的 context.Background()
  3. 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 一把梭。