用任务通道和结果通道把任务的提交和执行结果分离,这是 Go 并发编程中最工程化的写法 —— 生产者、worker、消费者各司其职。
📚 系列关系: 本文是 Worker Pool 的变体。与 信号量模式(每次新建协程)和 Worker Pool(固定 worker 复用)相比,双层 Channel 的核心 worker 复用机制完全一样,只是多了一条结果通道把计算结果传回调用方。
完整代码
1 | package main |
1 | package main |
1 | package main |
核心原理
双层 Channel:一条通道负责分发任务(jobs),另一条通道负责收集结果(res)。
worker 在中间充当桥梁:从 jobs 取任务,处理后把结果推入 res。
与 Worker Pool 的区别
| 模式 | 任务传递 | 结果获取 | 特点 |
|---|---|---|---|
| Worker Pool | jobs <- func() | 闭包内部处理 | 任务本身就是函数,结果在闭包内消化 |
| 双层 Channel | jobs <- data | res <- 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 | 主协程 worker 协程(3个) |
图解
双层 Channel 数据流
1 | jobs 通道(任务流入) |
与 Worker Pool 裸写版对比
把两篇的裸写版放一起,核心差异一目了然:
1 | // Worker Pool:1 条通道 + WaitGroup |
1 | // 双层 Channel:2 条通道 + 无 WaitGroup |
对比表
| 维度 | Worker Pool | 双层 Channel |
|---|---|---|
| 通道数量 | 1 条 jobs | 2 条 jobs + res |
| 同步机制 | sync.WaitGroup(3 个方法调用) | 无,靠 <-res 阻塞等结果 |
| 提交内容 | func(){...} 闭包 | int 纯数据 |
| worker 写法 | 内联匿名函数,闭包抓变量 | 独立函数,参数传递 |
| 结果处理 | 闭包内部直接处理 | 必须 main 里手动 <-res 收集 |
为什么双层 Channel 看起来更简单?
- 没有
WaitGroup,少了Add/Done/Wait3 个调用 - 不需要每次提交时包一层
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 | type Result struct { |
结果自带 ID,取出来就知道对应哪个任务。
方式 2:固定长度切片 + 索引
1 | results := make([]int, 10) // 预分配 |
worker 直接按索引写入切片,无需包装。注意并发写入时,每个 worker 写的位置互不冲突,不需要额外锁。
方式 3:顺序不重要
1 | // 当前文章写法 |
只要 10 个结果都拿到就行,适合批量处理场景(如并发下载图片,只要全部下载完就行)。
适用场景一览
| 场景 | 推荐方式 |
|---|---|
| 需要结果和输入一一对应 | 结构体带 ID |
| 结果写回固定位置 | 切片+索引 |
| 批量处理、汇总即可 | 直接取 |
| 闭包内部处理 | Worker Pool |
常见问题
结果的顺序为什么不一定?
多个 worker 并发执行,谁先处理完任务就把结果写入 res 通道,顺序取决于执行速度。如果你需要结果与输入一一对应,可以用上一节的 方式 1:结构体包装,把原始 ID 带回来。
如果任务数量不确定,怎么收集结果?
用 for range 监听 res 通道,但需要配合 sync.WaitGroup 来确定何时关闭 res:
1 | package main |
为什么要用单独的 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 处阻塞,无法自然退出。
解决方法:
- 确保消费者取完所有结果(如本文示例)
- 或者用
select+ 超时避免永久阻塞