在Go语言中做并发编程,最容易犯的错误就是拿到任务就无脑开goroutine。真正合理的做法,是先分析任务本身是否可独立、是否存在共享状态、数据量是否适合一次性并发。任务拆解的本质,是把一个粗粒度的目标,切成多个可调度、可控制、可恢复的细粒度工作单元。

一、为什么不能直接全量开goroutine
假设我们需要处理十万条用户数据,每条数据要调用一次外部接口并写库。初学者常写出下面这种代码:
package main
import (
"fmt"
"time"
)
func handle(id int) {
time.Sleep(10 * time.Millisecond)
fmt.Println("handled", id)
}
func main() {
for i := 0; i < 100000; i++ {
go handle(i)
}
time.Sleep(5 * time.Second)
}
这段代码在数量小时看不出问题,但到了十万级别,瞬间创建十万个goroutine会让调度器压力剧增,同时外部接口可能被打挂。Go的goroutine虽然轻量,但并非没有成本,每个 goroutine 默认有栈空间,大量阻塞在IO上也会占用调度资源。
更关键的是,这种写法完全失去了并发控制能力。你无法知道任务什么时候全部完成,也无法在出错时统一取消。因此任务拆解的第一步,永远是加一层并发度控制,而不是盲目并发。
二、按数据分片拆解任务
数据分片是最直观的拆解方式:把大数据集切成分块,每块交给一个worker处理。下面用sync.WaitGroup配合固定数量的goroutine实现:
package main
import (
"fmt"
"sync"
"time"
)
func worker(id int, jobs <-chan int, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs {
time.Sleep(10 * time.Millisecond)
fmt.Printf("worker %d handled %dn", id, job)
}
}
func main() {
total := 1000
workerCount := 10
jobs := make(chan int, 100)
var wg sync.WaitGroup
for i := 1; i <= workerCount; i++ {
wg.Add(1)
go worker(i, jobs, &wg)
}
for i := 0; i < total; i++ {
jobs <- i
}
close(jobs)
wg.Wait()
fmt.Println("all done")
}
这里我们把一千个任务通过带缓冲的channel分发给十个常驻worker。这种结构被称为worker pool,它的好处是并发数恒定,不会因为任务多就爆炸。jobs channel相当于一个任务队列,worker从里面取活儿,天然实现了负载均衡。
分片大小也值得斟酌。如果切片太大,单个worker处理太久会导致其他worker饿死;太小则channel收发开销占比上升。一般建议每片处理时间在毫秒到百毫秒级,并结合压测调整worker数量和channel缓冲。
三、按阶段扇出扇入拆解
当任务有明确阶段,比如先抓取、再解析、最后存储,可以用扇出扇入(fan-out fan-in)模式。扇出是多个goroutine消费同一输入,扇入是把多个结果合并。
package main
import (
"fmt"
"sync"
)
func gen(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
out <- n * n
}
close(out)
}()
return out
}
func merge(chans ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
output := func(c <-chan int) {
defer wg.Done()
for n := range c {
out <- n
}
}
wg.Add(len(chans))
for _, c := range chans {
go output(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
in := gen(1, 2, 3, 4)
c1 := square(in)
c2 := square(in)
for n := range merge(c1, c2) {
fmt.Println(n)
}
}
上面的例子里,两个square goroutine同时处理gen产出的数据,这就是扇出;merge把c1和c2的结果收拢,就是扇入。这种拆解让CPU密集型和IO密集型阶段可以重叠执行,整体吞吐比串行高得多。
要注意的是,扇入阶段的merge必须等所有上游关闭才能关输出channel,否则会向已关闭channel写数据引发panic。用WaitGroup跟踪上游是标准做法。
四、用context控制取消与超时
拆解后的任务若某个环节失败,应该能通知所有goroutine停下来,而不是继续空耗。context包就是干这个的。
package main
import (
"context"
"fmt"
"time"
)
func worker(ctx context.Context, id int) {
for {
select {
case <-ctx.Done():
fmt.Printf("worker %d cancelled: %vn", id, ctx.Err())
return
default:
time.Sleep(100 * time.Millisecond)
fmt.Printf("worker %d workingn", id)
}
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancel()
for i := 1; i <= 3; i++ {
go worker(ctx, i)
}
<-ctx.Done()
fmt.Println("main exit")
}
通过WithTimeout创建的ctx,在到达时限后会让Done channel可读,所有监听它的worker立刻退出。如果是错误驱动取消,调用cancel()函数即可。把ctx作为参数透传到每个拆解出的子任务,是整个并发体系可控的关键。
实际项目中,往往把context和worker pool结合:主goroutine收到取消信号,关闭任务channel或借助ctx让worker在取任务前就退出,避免已发出请求的资源浪费。
五、拆解时的常见误区
第一个误区是过度拆解。把本来只需微秒级的操作拆成多goroutine,调度和channel通信的损耗反而盖过了并发收益。第二个误区是共享变量不保护,多个拆解任务写同一map却不加锁,导致竞态。
| 误区 | 表现 | 正确做法 |
|---|---|---|
| 无限制goroutine | 内存暴涨、接口限流 | worker pool限并发 |
| 过细拆解 | 吞吐反降 | 按耗时评估粒度 |
| 忽略取消 | 僵尸goroutine | 传context控制 |
任务拆解不是银弹,它要求你理解数据流向和资源边界。写之前在纸上画一下阶段和依赖,比直接敲代码更省时间。
六、总结实战思路
面对一个Go并发需求,先问三个问题:任务能否独立?瓶颈在CPU还是IO?失败要不要整体停?回答清楚后,选数据分片还是阶段扇出,再加worker数和context,基本就能写出稳的并发程序。
拆解的核心价值,是把不可控的洪水式并发,变成可观测、可调度、可取消的工作流。多写几遍worker pool和扇出扇入,你会对Go的并发模型有实感。