任务通道结果通道把任务的提交和执行结果分离,这是 Go 并发编程中最工程化的写法 —— 生产者、worker、消费者各司其职。

📚 系列关系: 本文是 Worker Pool 的变体。与 信号量模式(每次新建协程)和 Worker Pool(固定 worker 复用)相比,双层 Channel 的核心 worker 复用机制完全一样,只是多了一条结果通道把计算结果传回调用方。


完整代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
package main

import (
"fmt"
"time"
)

// worker 从 jobs 通道取任务,把结果写入 res 通道
func worker(id int, jobs <-chan int, res chan<- int) {
for j := range jobs {
time.Sleep(time.Second) // 模拟耗时操作
res <- j * 2 // 结果写入结果通道
}
}

func main() {
jobs := make(chan int, 3) // 任务通道(只进)
res := make(chan int, 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)
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
package main

import (
"fmt"
"time"
)

func worker(id int, jobs <-chan int, res chan<- int) {
for j := range jobs {
time.Sleep(time.Second)
res <- j * 2
}
}

func main() {
jobs := make(chan int, 3)
res := make(chan int, 3)

for i := 0; i < 3; i++ {
go worker(i, jobs, res)
}

for i := 0; i < 10; i++ {
jobs <- i
}
close(jobs)

for i := 0; i < 10; i++ {
fmt.Println(<-res)
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
package main

import (
"fmt"
"sync"
"time"
)

func worker(id int, jobs <-chan int, res chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for j := range jobs {
time.Sleep(time.Second)
res <- j * 2
}
}

func main() {
jobs := make(chan int, 3)
res := make(chan int, 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)

go func() {
wg.Wait()
close(res)
}()

for r := range res {
fmt.Println(r)
}
}

核心原理

双层 Channel:一条通道负责分发任务(jobs),另一条通道负责收集结果(res)。

worker 在中间充当桥梁:从 jobs 取任务,处理后把结果推入 res。

与 Worker Pool 的区别

模式任务传递结果获取特点
Worker Pooljobs <- func()闭包内部处理任务本身就是函数,结果在闭包内消化
双层 Channeljobs <- datares <- result任务和数据分离,worker 处理完写到 res 通道

关键语法:通道方向

func 参数中的 <-chan 和 chan<-

1
func worker(id int, jobs <-chan int, res chan<- int) {

Go 支持通道方向的语法限制,让函数签名更清晰地表达意图:

语法含义操作
<-chan int只读通道只能 <-jobs(从通道取数据)
chan<- int只写通道只能 res <- x(往通道发数据)
chan int双向通道既能发也能收

好处: 编译器会阻止你在 worker 中错误地使用 close(jobs)close(res),从语法层面杜绝误操作。


执行流程

步骤 1 — 创建两个通道

jobs(任务通道)和 res(结果通道),缓冲都为 3。

步骤 2 — 启动 worker

3 个 worker 同时监听 jobs 通道,处理完把结果推入 res 通道。

步骤 3 — 提交任务

主协程往 jobs 通道写入 10 个任务,worker 自动取出执行。

步骤 4 — 关闭任务通道

close(jobs) 通知所有 worker 没有更多任务了,for range 循环结束。

步骤 5 — 收集结果

主协程从 res 通道读取 10 个结果,顺序不一定等于提交顺序。


执行时序图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
主协程                            worker 协程(3个)

├─ jobs = make(chan int, 3)
├─ res = make(chan int, 3)

├─ go worker(0, jobs, res) ────────────────→ for range jobs 等待...
├─ go worker(1, jobs, res) ────────────────→ for range jobs 等待...
├─ go worker(2, jobs, res) ────────────────→ for range jobs 等待...

├─ jobs <- 0 ────────────────────────────→ worker0 取到,计算 0*2=0,res <- 0
├─ jobs <- 1 ────────────────────────────→ worker1 取到,计算 1*2=2,res <- 2
├─ jobs <- 2 ────────────────────────────→ worker2 取到,计算 2*2=4,res <- 4
│ ...后续任务依此类推...

├─ close(jobs)
│ worker 的 for range 结束,协程退出

├─ <-res (第1次) → 0
├─ <-res (第2次) → 2
├─ <-res (第3次) → 4
├─ ...直到10个结果全部取出

└─ 程序退出

图解

双层 Channel 数据流

1
2
3
4
5
6
7
8
9
10
11
            jobs 通道(任务流入)

┌────────────┼────────────┐
↓ ↓ ↓
worker0 worker1 worker2
│ │ │
└────────────┼────────────┘

res 通道(结果流出)

fmt.Println(<-res)

与 Worker Pool 裸写版对比

把两篇的裸写版放一起,核心差异一目了然:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// Worker Pool:1 条通道 + WaitGroup
jobs := make(chan func(), 3)
var wg sync.WaitGroup

for i := 0; i < 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for j := range jobs {
j() // 执行闭包函数
}
}(i)
}

for i := 0; i < 10; i++ {
jobs <- func() { /* 业务逻辑 */ } // 提交的是闭包
}
close(jobs)
wg.Wait()
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 双层 Channel:2 条通道 + 无 WaitGroup
jobs := make(chan int, 3)
res := make(chan int, 3)

for i := 0; i < 3; i++ {
go worker(i, jobs, res)
}

for i := 0; i < 10; i++ {
jobs <- i // 提交的是纯数据
}
close(jobs)

for i := 0; i < 10; i++ {
fmt.Println(<-res) // 手动拉取结果
}

对比表

维度Worker Pool双层 Channel
通道数量1 条 jobs2 条 jobs + res
同步机制sync.WaitGroup(3 个方法调用)无,靠 <-res 阻塞等结果
提交内容func(){...} 闭包int 纯数据
worker 写法内联匿名函数,闭包抓变量独立函数,参数传递
结果处理闭包内部直接处理必须 main 里手动 <-res 收集

为什么双层 Channel 看起来更简单?

  • 没有 WaitGroup,少了 Add/Done/Wait 3 个调用
  • 不需要每次提交时包一层 func(){...}
  • worker 拆成独立函数,main 函数更干净

但双层 Channel 有隐藏复杂度:

  • res 缓冲区满了会阻塞 worker — 消费者没取完,worker 卡在 res <- j*2
  • 结果数量必须提前知道(写死 for i := 0; i < 10),否则要配合 WaitGroup 来关 res
  • 结果和输入不一定对应,需要自己解决顺序问题(见下一节)

Worker Pool 闭包内部消化结果不用管,但闭包语法和 WaitGroup 增加了理解门槛。双层 Channel 的”简单”在于少了样板代码,但代价是结果收集的责任从 worker 转移到了调用方。


结果与任务的对应问题

核心问题: 提交顺序是 0, 1, 2, 3, 4...,取出顺序可能是 2, 0, 4, 1, 3...。因为多个 worker 并发执行,谁先处理完谁先写入 res 通道。

Worker Pool 没有这个问题——闭包里同时有输入和输出,天然绑定在一起。双层 Channel 把结果和输入分离了,所以需要额外处理。

方式 1:结果包装结构体(推荐)

1
2
3
4
5
6
7
8
9
10
11
12
13
type Result struct {
ID int
Data int
}

res := make(chan Result, 3)

// worker 内部:
res <- Result{ID: j, Data: j * 2}

// main 中:
r := <-res
fmt.Printf("任务%d → 结果%d\n", r.ID, r.Data)

结果自带 ID,取出来就知道对应哪个任务。

方式 2:固定长度切片 + 索引

1
2
3
4
5
6
7
8
9
results := make([]int, 10)  // 预分配

// worker 内部:
results[j] = j * 2 // 直接写到对应位置

// main 中:
for _, r := range results {
fmt.Println(r)
}

worker 直接按索引写入切片,无需包装。注意并发写入时,每个 worker 写的位置互不冲突,不需要额外锁。

方式 3:顺序不重要

1
2
3
4
// 当前文章写法
for i := 0; i < 10; i++ {
fmt.Println(<-res) // 不关心对应关系
}

只要 10 个结果都拿到就行,适合批量处理场景(如并发下载图片,只要全部下载完就行)。

适用场景一览

场景推荐方式
需要结果和输入一一对应结构体带 ID
结果写回固定位置切片+索引
批量处理、汇总即可直接取
闭包内部处理Worker Pool

常见问题

结果的顺序为什么不一定?

多个 worker 并发执行,谁先处理完任务就把结果写入 res 通道,顺序取决于执行速度。如果你需要结果与输入一一对应,可以用上一节的 方式 1:结构体包装,把原始 ID 带回来。

如果任务数量不确定,怎么收集结果?

for range 监听 res 通道,但需要配合 sync.WaitGroup 来确定何时关闭 res:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
package main

import (
"fmt"
"sync"
"time"
)

func worker(id int, jobs <-chan int, res chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for j := range jobs {
time.Sleep(time.Second)
res <- j * 2
}
}

func main() {
jobs := make(chan int, 3)
res := make(chan int, 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
go func() {
wg.Wait()
close(res)
}()

// 安全读取
for r := range res {
fmt.Println(r)
}
}

为什么要用单独的 worker 函数,而不是 main 里直接写 go func()?

分离 worker 函数有两个好处:

  1. 复用性 — 多个地方可以创建相同逻辑的 worker
  2. 可测试性 — 单独测试 worker 的逻辑,不需要写完整的 main

当然脚本场景直接 go func() 也没问题,看代码组织需求。

双层 Channel 的 worker 怎么优雅退出?

close(jobs) 后 worker 的 for range jobs 会自然结束,协程退出。但要注意:如果 res 通道的消费者(main 中 <-res)没有全部取出结果,res 缓冲区满后 worker 会在 res <- j*2 处阻塞,无法自然退出。

解决方法:

  • 确保消费者取完所有结果(如本文示例)
  • 或者用 select + 超时避免永久阻塞