// worker 从 jobs 通道取任务,把结果写入 res 通道 funcworker(id int, jobs <-chanint, res chan<- int) { for j := range jobs { time.Sleep(time.Second) // 模拟耗时操作 res <- j * 2// 结果写入结果通道 } }
funcmain() { jobs := make(chanint, 3) // 任务通道(只进) res := make(chanint, 3) // 结果通道(只出)
// 启动 3 个 worker for i := 0; i < 3; i++ { go worker(i, jobs, res) }
// 提交 10 个任务 for i := 0; i < 10; i++ { jobs <- i } close(jobs) // 通知 worker 没有更多任务了
// 收集 10 个结果 for i := 0; i < 10; i++ { fmt.Println(<-res) } }
funcworker(id int, jobs <-chanint, res chan<- int, wg *sync.WaitGroup) { defer wg.Done() for j := range jobs { time.Sleep(time.Second) res <- j * 2 } }
funcmain() { jobs := make(chanint, 3) res := make(chanint, 3) var wg sync.WaitGroup
for i := 0; i < 3; i++ { wg.Add(1) go worker(i, jobs, res, &wg) }
for i := 0; i < 10; i++ { jobs <- i } close(jobs)
gofunc() { wg.Wait() close(res) }()
for r := range res { fmt.Println(r) } }
核心原理
双层 Channel:一条通道负责分发任务(jobs),另一条通道负责收集结果(res)。
worker 在中间充当桥梁:从 jobs 取任务,处理后把结果推入 res。
与 Worker Pool 的区别
模式
任务传递
结果获取
特点
Worker Pool
jobs <- func()
闭包内部处理
任务本身就是函数,结果在闭包内消化
双层 Channel
jobs <- data
res <- result
任务和数据分离,worker 处理完写到 res 通道
关键语法:通道方向
func 参数中的 <-chan 和 chan<-
1
funcworker(id int, jobs <-chanint, res chan<- int) {
funcworker(id int, jobs <-chanint, res chan<- int, wg *sync.WaitGroup) { defer wg.Done() for j := range jobs { time.Sleep(time.Second) res <- j * 2 } }
funcmain() { jobs := make(chanint, 3) res := make(chanint, 3) var wg sync.WaitGroup
// 启动 worker for i := 0; i < 3; i++ { wg.Add(1) go worker(i, jobs, res, &wg) }
// 提交任务 for i := 0; i < 10; i++ { jobs <- i } close(jobs)
// 等待所有 worker 完成,再关闭 res gofunc() { wg.Wait() close(res) }()
// 安全读取 for r := range res { fmt.Println(r) } }
为什么要用单独的 worker 函数,而不是 main 里直接写 go func()?
分离 worker 函数有两个好处:
复用性 — 多个地方可以创建相同逻辑的 worker
可测试性 — 单独测试 worker 的逻辑,不需要写完整的 main
当然脚本场景直接 go func() 也没问题,看代码组织需求。
双层 Channel 的 worker 怎么优雅退出?
close(jobs) 后 worker 的 for range jobs 会自然结束,协程退出。但要注意:如果 res 通道的消费者(main 中 <-res)没有全部取出结果,res 缓冲区满后 worker 会在 res <- j*2 处阻塞,无法自然退出。