如何在发生错误时停止从Go Channel读取数据?
使用errgroup.WithContext实现错误时立即停止Channel读取
原代码中,当某个goroutine返回错误后,其他goroutine仍会继续从channel读取数据直到channel关闭,无法实现遇到错误立即停止的需求。通过errgroup.WithContext可以在任意goroutine出错时触发上下文取消,让所有goroutine同步停止工作,具体整合方式如下:
修改后的完整代码
package main import ( "context" "fmt" "time" "golang.org/x/sync/errgroup" ) func main() { const threads = 3 ch := make(chan int, threads) // 创建带上下文的errgroup,任意goroutine返回错误时,ctx会被取消 eg, ctx := errgroup.WithContext(context.Background()) for i := 0; i < threads; i++ { i := i eg.Go(func() error { fmt.Printf("Thread %d: STARTED\n", i) for { select { // 监听上下文取消信号,一旦触发就退出循环 case <-ctx.Done(): fmt.Printf("Thread %d: STOPPED due to context cancel\n", i) return ctx.Err() // 读取channel数据 case n, ok := <-ch: if !ok { // channel关闭,正常退出 return nil } fmt.Printf("Thread %d: GOT=%d\n", i, n) time.Sleep(time.Duration(1) * time.Second) // 模拟线程失败 if n == 2 { fmt.Printf("Thread %d: FAILED\n", i) return fmt.Errorf("Thread %d: FAILED", i) } } } }) } // 发送数据到channel,同时监听上下文取消,避免阻塞泄漏 go func() { defer close(ch) for i := 0; i < 9; i++ { select { case <-ctx.Done(): fmt.Printf("Sender: STOPPED due to context cancel\n") return case ch <- i: } } }() if err := eg.Wait(); err != nil { panic(err) } }
关键改动说明
- 绑定上下文与errgroup:使用
errgroup.WithContext创建group,替代原有的无上下文group。当任意goroutine返回错误时,该group会自动取消绑定的上下文,所有监听该上下文的goroutine都会收到取消信号。 - goroutine内监听取消信号:将原有的
for n := range ch改为select语句,同时监听channel读取和ctx.Done()信号。一旦上下文被取消,立即退出循环并停止处理。 - 发送端监听取消信号:单独开一个goroutine发送数据,并监听上下文取消,避免当所有消费goroutine都退出后,发送操作因channel满而阻塞,造成goroutine泄漏。
预期运行输出
Thread 2: STARTED Thread 1: STARTED Thread 0: STARTED Thread 2: GOT=0 Thread 1: GOT=1 Thread 0: GOT=2 Thread 0: FAILED Thread 1: STOPPED due to context cancel Thread 2: STOPPED due to context cancel Sender: STOPPED due to context cancel panic: Thread 0: FAILED
可以看到,当Thread 0失败后,其他线程和发送端都会立即停止工作,符合“遇到错误立即停止读取”的需求。
内容的提问来源于stack exchange,提问作者TheAschr
相关产品推荐
相关产品推荐

