Go语言:从并发goroutine收集响应仅获半数结果的问题排查
问题:并行任务响应读取异常仅获取半数结果
我需要基于输入数组运行并行任务,等待所有任务完成后处理它们的响应。我的代码使用WaitGroup等待所有goroutine执行完毕,再从通道读取响应,但仅能获取到半数响应。请问这种获取goroutine响应的方式是否正确?若正确我遗漏了什么?若不正确,正确的实现方式是什么?
简化后的代码:
package main import ( "fmt" "sync" "time" ) func odd(i int) (int, error) { time.Sleep(1 * time.Second) if i%2 == 0 { return i, fmt.Errorf("even number") } else { return i, nil } } func main() { type R struct { val int err error } wg := sync.WaitGroup{} respChan := make(chan R, 10) for i := 0; i < 10; i++ { wg.Add(1) go func(i int) { defer wg.Done() val, err := odd(i) r := R{val: val, err: err} respChan <- r fmt.Printf("Queued Response: %d , size: %d \n", r.val, len(respChan)) }(i) } wg.Wait() fmt.Println("Done Waiting") fmt.Println("Response Channel Length: ", len(respChan)) for i := 0; i < len(respChan); i++ { r := <-respChan if r.err != nil { fmt.Printf("[%d] : %d , %s\n", i, r.val, r.err.Error()) } else { fmt.Printf("[%d] : %d\n", i, r.val) } } fmt.Println("Finished") }
输出结果:
Queued Response: 5 , size: 1 Queued Response: 0 , size: 2 Queued Response: 2 , size: 3 Queued Response: 9 , size: 7 Queued Response: 3 , size: 5 Queued Response: 4 , size: 6 Queued Response: 1 , size: 4 Queued Response: 7 , size: 8 Queued Response: 6 , size: 9 Queued Response: 8 , size: 10 Done Waiting Response Channel Length: 10 [0] : 5 [1] : 0 , even number [2] : 2 , even number [3] : 1 [4] : 3 Finished
错误原因
核心问题出在读取通道的循环条件:for i := 0; i < len(respChan); i++。
每次从通道读取元素时,len(respChan)会实时递减:
- 初始时通道长度是10,i=0,满足0<10,读取第一个元素后通道长度变为9,i自增到1
- 第二次循环:i=1<9,读取第二个元素后通道长度变为8,i自增到2
- ...以此类推,当i=5时,通道长度已递减到5,此时5<5不成立,循环直接终止,最终只读取了5个元素。
这种依赖通道实时长度作为循环终止条件的方式完全错误,因为通道长度会随着读取操作动态变化。
正确实现方式
有两种可靠的实现方式:
方式1:使用固定次数循环(已知任务总数)
既然明确知道要处理10个任务的响应,直接循环10次即可,不受通道长度变化影响:
修改main函数中的读取循环部分:
fmt.Println("Done Waiting") fmt.Println("Response Channel Length: ", len(respChan)) // 固定循环10次,对应10个任务的响应 for i := 0; i < 10; i++ { r := <-respChan if r.err != nil { fmt.Printf("[%d] : %d , %s\n", i, r.val, r.err.Error()) } else { fmt.Printf("[%d] : %d\n", i, r.val) } }
方式2:关闭通道后用range循环读取(更通用)
在所有goroutine完成后关闭通道,然后用range遍历通道,这样会自动读取所有元素直到通道关闭,无需关心具体数量:
修改main函数:
wg.Wait() close(respChan) // 所有goroutine完成后关闭通道 fmt.Println("Done Waiting") fmt.Println("Response Channel Length: ", len(respChan)) i := 0 // range会读取通道中所有元素,直到通道关闭 for r := range respChan { if r.err != nil { fmt.Printf("[%d] : %d , %s\n", i, r.val, r.err.Error()) } else { fmt.Printf("[%d] : %d\n", i, r.val) } i++ }
这种方式更通用,尤其适用于任务数量不确定的场景,且符合Go语言通道的使用规范。
内容的提问来源于stack exchange,提问作者Viral
相关产品推荐
相关产品推荐

