Go channel 流水线入门:生产、处理、取消和关闭怎么配合

用批量处理任务讲 Go channel 流水线的基础写法,包括生产者、工作者、结果汇总、context 取消和 channel 关闭规则,以及多级流水线、背压控制和超时处理。

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)
}

这时 produceworker 都会因为 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 Jobresults 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 帮助排查:

  1. 每个 goroutine 都有退出路径:检查每个 go func() 是否能通过 returnclose(channel) 结束
  2. channel 关闭后 range 会退出:确保发送方会关闭 channel
  3. context 取消要传播到每一层:从入口到最底层,每层都 select ctx.Done()
  4. 不要用无限制 goroutine:worker 数量应该是固定的,不要每来一个任务就新建 goroutine
  5. 阻塞发送要有超时或取消:裸写 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 流水线在生产环境运行时,还需要关注:

  1. 监控指标:暴露正在运行的 worker 数量、channel 长度、处理速率等指标
  2. 优雅关闭:服务停止时,先取消 context 等待所有 goroutine 退出,再关闭资源
  3. 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 当成常规错误处理。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 熔断、降级与限流:Go 微服务韧性设计完全指南
  2. 事件溯源与 CQRS 在 Go 中的实践:复杂业务系统的架构升级
  3. TinyGo 嵌入式开发与物联网实战:微控制器编程完全指南