channel 是 Go 很有代表性的特性,但很多初学者会把它写成“能跑就行”的并发代码:谁关闭 channel 不清楚,错误怎么返回不清楚,取消后 goroutine 会不会退出也不清楚。短期看只是偶尔卡住,长期看就是 goroutine 泄漏和线上不稳定。
本文用一个批量处理 URL 的例子,讲一个简单流水线:生产任务,多个 worker 处理,最后汇总结果。重点不是追求最高性能,而是把 channel 关闭、context 取消和错误处理讲清楚。
定义任务和结果
先定义结构:
type Job struct {
ID int
URL string
}
type Result struct {
JobID int
Code int
Err error
}
任务包含 URL,结果包含状态码或错误。不要只传字符串。真实项目里,多一点结构字段会让日志和排查更清楚。
生产者负责关闭任务 channel
生产任务:
func produce(ctx context.Context, urls []string) <-chan Job {
jobs := make(chan Job)
go func() {
defer close(jobs)
for i, u := range urls {
select {
case jobs <- Job{ID: i + 1, URL: u}:
case <-ctx.Done():
return
}
}
}()
return jobs
}
谁发送,谁关闭。produce 是任务 channel 的唯一发送方,所以由它关闭 jobs。接收方不要关闭别人发送的 channel。这个规则很简单,却能避免很多 panic。
select 里监听 ctx.Done(),表示如果上层取消,生产者会停止发送并退出。
worker 处理任务
worker 从 jobs 读取,向 results 写入:
func worker(ctx context.Context, client *http.Client, jobs <-chan Job, results chan<- Result) {
for job := range jobs {
code, err := fetch(ctx, client, job.URL)
select {
case results <- Result{JobID: job.ID, Code: code, Err: err}:
case <-ctx.Done():
return
}
}
}
注意 channel 方向:jobs <-chan Job 表示只读,results chan<- Result 表示只写。方向不是必须写,但写上后函数意图更清楚,编译器也能帮你挡住误用。
fetch 也要使用 context:
func fetch(ctx context.Context, client *http.Client, rawURL string) (int, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, rawURL, nil)
if err != nil {
return 0, err
}
resp, err := client.Do(req)
if err != nil {
return 0, err
}
defer resp.Body.Close()
io.Copy(io.Discard, resp.Body)
return resp.StatusCode, nil
}
这里把响应体读完并关闭,是为了连接复用。即使不关心 body,也不要直接丢掉。
汇总并关闭结果 channel
启动多个 worker:
func run(ctx context.Context, urls []string, workerCount int) []Result {
client := &http.Client{Timeout: 5 * time.Second}
jobs := produce(ctx, urls)
results := make(chan Result)
var wg sync.WaitGroup
for i := 0; i < workerCount; i++ {
wg.Add(1)
go func() {
defer wg.Done()
worker(ctx, client, jobs, results)
}()
}
go func() {
wg.Wait()
close(results)
}()
var out []Result
for result := range results {
out = append(out, result)
}
return out
}
结果 channel 由谁关闭?不是 worker 单独关闭,因为有多个 worker,任何一个提前关闭都会让其他 worker 发送时 panic。正确做法是等待所有 worker 结束后,由一个单独 goroutine 关闭 results。
这个模式很常见:多个发送者共享一个输出 channel 时,用 WaitGroup 等所有发送者结束,再关闭输出。
取消后不要卡住
如果上层只想要前 10 个成功结果,可以在拿够后取消 context:
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
results := run(ctx, urls, 5)
_ = results
在更复杂的版本里,你可能边读结果边决定取消。关键是所有发送和接收都要能响应 ctx.Done(),否则某个 goroutine 可能永远卡在发送上。
一个危险写法是:
results <- result
如果没人接收,这行会一直阻塞。更稳的写法是前面 worker 中的 select。channel 流水线里,取消路径和正常路径一样重要。
错误处理策略
遇到单个 URL 失败,要不要取消全部任务?这取决于业务。如果是批量探测网站,某个失败可以记录结果继续;如果是多步骤导入,其中一个失败就必须停止,那就应该取消 context。
可以把取消权放在汇总层:
for result := range results {
if result.Err != nil {
cancel()
return nil, result.Err
}
out = append(out, result)
}
这时 produce 和 worker 都会因为 context 取消而退出。错误策略不要藏在 worker 里,否则后面很难理解为什么任务突然停了。
buffer 要谨慎
给 channel 加 buffer 可以改善吞吐,但不要把 buffer 当成修复死锁的工具:
results := make(chan Result, workerCount)
小 buffer 可以减少 worker 和汇总 goroutine 的等待。但如果程序逻辑依赖“buffer 足够大所以不会卡住”,那只是把问题推迟。正确的流水线应该在无 buffer 或小 buffer 下也能正常退出。
测试是否能退出
并发流水线最好写一个取消测试,确认取消后函数能返回:
func TestRunCanBeCanceled(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
done := make(chan struct{})
go func() {
_ = run(ctx, []string{"https://example.com"}, 2)
close(done)
}()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("pipeline did not stop after cancel")
}
}
这个测试不追求覆盖所有网络行为,它只验证生命周期。很多 channel bug 不是结果算错,而是程序结束不了。能用测试守住退出路径,并发代码会更可信。
小结
Go channel 流水线的基本规则是:发送方关闭 channel,多个发送方要用 WaitGroup 汇合后再关闭,所有可能阻塞的发送和接收都要考虑 context 取消。worker 不应该擅自关闭共享输出 channel,错误策略最好集中在汇总层。
channel 不是为了写炫酷并发,而是为了表达数据流。初学者先把生命周期写清楚,再考虑性能优化。能正常结束的并发程序,才有资格谈快。
多级流水线
本文的流水线只有"生产 -> 处理 -> 汇总"三层。更复杂的场景可能有多级处理,比如"生产 -> 解析 -> 验证 -> 处理 -> 汇总"。每级之间都是一个 channel。
func parse(ctx context.Context, jobs <-chan Job) <-chan ParsedJob {
out := make(chan ParsedJob)
go func() {
defer close(out)
for job := range jobs {
parsed, err := parseURL(job.URL)
select {
case out <- ParsedJob{Job: job, Parsed: parsed, Err: err}:
case <-ctx.Done():
return
}
}
}()
return out
}
每级的关闭规则都一样:谁发送谁关闭。读者只需要理解一级,多级只是把同一个模式串联起来。
背压控制
如果生产者太快、消费者太慢,中间 channel 的 buffer 会填满,然后生产者阻塞。这是天然的背压(backpressure)。但如果你的设计不允许阻塞,比如数据来自消息队列,过度堆积会撑爆内存。
一种简单策略是丢弃过旧数据:
select {
case jobs <- job:
default:
// buffer 满时丢弃
log.Printf("dropping job %d, buffer full", job.ID)
}
另一种策略是慢启动消费者,或者增加 worker。背压策略取决于业务对数据丢失的容忍度。
Pipeline 超时处理
单条任务不应该让整个流水线卡住。可以给每条任务设置独立超时:
func fetchWithTimeout(ctx context.Context, client *http.Client, url string) (int, error) {
ctx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
resp, err := client.Do(req)
if err != nil {
return 0, err
}
defer resp.Body.Close()
io.Copy(io.Discard, resp.Body)
return resp.StatusCode, nil
}
注意这里的 WithTimeout 是从传入的 ctx 派生,而不是 context.Background()。这样上层取消仍然能穿透。
常见 Pitfall:关闭已关闭的 channel
Go 里关闭已经关闭的 channel 会 panic。多人协作时容易踩这个坑:
// worker A 和 worker B 都尝试关闭 results
func worker(jobs <-chan Job, results chan<- Result) {
for job := range jobs {
// ...
}
close(results) // 致命错误!多个 worker 会同时尝试关闭
}
正确的模式前文已讲:用 WaitGroup 等所有发送者结束,再由一个单独的 goroutine 关闭。这是 channel 流水线最重要的模式之一。
channel 方向不是可选装饰
本文示例写了 jobs <-chan Job 和 results chan<- Result。有人觉得方向可以省略,代码更短。但实际上方向能拦住一类常见 bug:
// 如果 worker 签名是 func worker(jobs chan Job, results chan Result)
// 编译器不会阻止你在 worker 里读取 results 或写入 jobs
// 方向约束后,误用会在编译期报错
编译期的检查永远比运行期的 panic 便宜。别省略 channel 方向。
Pipeline 模式对比
| 模式 | 特点 | 适用场景 |
|---|---|---|
| 无缓冲 channel | 同步握手,严格顺序 | 需要确认每步完成的场景 |
| 有缓冲 channel | 异步解耦,减少等待 | 生产者和消费者速度不匹配 |
| 多级流水线 | 职责分离,可组合 | 复杂数据处理链路 |
| Fan-out/Fan-in | 并行处理,最后汇总 | 批量任务处理 |
本文的流水线本质上是 Fan-out(多个 worker)+ Fan-in(汇总结果)的组合。这是 Go 并发里最实用的模式之一。
FAQ
Q: 为什么不用 errgroup 代替 WaitGroup?
A: errgroup 是一个好工具,自带 error 收集和 context 取消。但它是一个额外依赖(golang.org/x/sync)。学习阶段建议先用 WaitGroup 理解原理,再引入 errgroup 简化代码。
Q: worker 数量怎么确定?
A: 主要取决于下游并发承受能力和 CPU 核心数。一个简单启发式:从 runtime.NumCPU() 开始,根据实测调整。IO 密集型任务可以适当增加,CPU 密集型不要太多。
Q: channel 可以跨 goroutine panic 恢复吗?
A: 不能也不应该。如果某个 goroutine panic,整个程序会崩溃。不要在 channel 操作里依赖 recover。正确做法是把错误编码到 Result 结构里传递。
小结
Go channel 流水线的基本规则是:发送方关闭 channel,多个发送方要用 WaitGroup 汇合后再关闭,所有可能阻塞的发送和接收都要考虑 context 取消。worker 不应该擅自关闭共享输出 channel,错误策略最好集中在汇总层。
channel 不是为了写炫酷并发,而是为了表达数据流。初学者先把生命周期写清楚,再考虑性能优化。能正常结束的并发程序,才有资格谈快。
Pipeline 中的错误收集
前面的示例把 error 放在 Result 里传递,这是最常见的做法。但如果需要收集所有错误而不是遇到第一个就停,可以换一种写法:
type Result struct {
JobID int
Code int
Err error
}
func collectResults(results <-chan Result) ([]Result, []error) {
var successes []Result
var errs []error
for r := range results {
if r.Err != nil {
errs = append(errs, fmt.Errorf("job %d: %w", r.JobID, r.Err))
} else {
successes = append(successes, r)
}
}
return successes, errs
}
这样上层可以决定怎么处理错误:全部返回、记录日志、还是重试。不要把"遇到错误就 panic"写在 goroutine 里。
防止 goroutine 泄漏 checklist
channel 流水线容易泄漏 goroutine,写一个 checklist 帮助排查:
- 每个 goroutine 都有退出路径:检查每个
go func()是否能通过return或close(channel)结束 - channel 关闭后 range 会退出:确保发送方会关闭 channel
- context 取消要传播到每一层:从入口到最底层,每层都 select
ctx.Done() - 不要用无限制 goroutine:worker 数量应该是固定的,不要每来一个任务就新建 goroutine
- 阻塞发送要有超时或取消:裸写
ch <- val是危险信号,优先用 select + ctx.Done()
可以用 go.uber.org/goleak 在测试中检测 goroutine 泄漏:
func TestMain(m *testing.M) {
goleak.VerifyTestMain(m)
}
带优先级的 Pipeline
有些任务更紧急,需要优先处理。可以用多个 channel 表达优先级:
func worker(jobs <-chan Job, urgentJobs <-chan Job, results chan<- Result) {
for {
select {
case job := <-urgentJobs:
process(job, results)
default:
select {
case job := <-urgentJobs:
process(job, results)
case job := <-jobs:
process(job, results)
}
}
}
}
外层 select 优先读取 urgentJobs,只有当没有紧急任务时才 fallback 到普通任务。这种写法避免了饥饿问题。
生产环境注意事项
channel 流水线在生产环境运行时,还需要关注:
- 监控指标:暴露正在运行的 worker 数量、channel 长度、处理速率等指标
- 优雅关闭:服务停止时,先取消 context 等待所有 goroutine 退出,再关闭资源
- panic 恢复:最外层的 goroutine 可以用
recover防止单个 worker panic 拖垮整个程序
func safeWorker(ctx context.Context, jobs <-chan Job, results chan<- Result) {
defer func() {
if r := recover(); r != nil {
log.Printf("worker panic: %v", r)
}
}()
worker(ctx, jobs, results)
}
但注意 panic 恢复只在最外层使用,不要把 recover 当成常规错误处理。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。