Go 并发过滤流水线中如何丢弃无效值且仅使用单输出通道?
Go过滤型并发流水线实现问题解答
问题1:移除非法值通道、仅保留合法值输出的修改方法
直接删除所有invalidValues相关的逻辑即可,不需要保留这个通道,消费者校验出非法值后直接跳过,仅对合法值执行发送操作:
修改后的消费者逻辑
//Consumers var wg sync.WaitGroup wg.Add(Workers) for i := 0; i < Workers; i++ { go func() { for value := range inputStream { // 校验不通过直接跳过 if !checkValue(value) { time.Sleep(5 * time.Second) continue } // 仅合法值走发送逻辑 select { case outputStream <- value: case <-done: return } time.Sleep(5 * time.Second) } wg.Done() }() }
其余配套修改:
- 删掉
invalidValues通道的定义 - 消费者结束后的关闭逻辑只需要关闭
outputStream - 删掉
errorFile相关的创建、写入逻辑,wg2只需要加1,仅保留写入合法结果的goroutine即可
问题2:select放入校验分支的修复方案
你之前尝试失败的核心原因是非法值分支不需要阻塞等待done通道,校验不通过时直接跳过当前值处理下一个输入即可,只有发送合法值时需要用select同时处理发送逻辑和全局退出信号,避免通道阻塞时程序无法退出。
示例正确写法:
for value := range inputStream { dataToWrite := value if valid := checkValue(value); valid { // 合法分支走发送+退出监听逻辑 select { case outputStream <- dataToWrite: case <-done: return } } // 非法分支什么都不用做,直接进入下一轮循环即可 time.Sleep(5 * time.Second) }
问题3:移除invalidValues读取goroutine后程序无法退出的原因及解决方案
根因
你定义的invalidValues是无缓冲通道,无缓冲通道的发送操作会一直阻塞,直到有其他goroutine执行接收操作。移除读取invalidValues的goroutine后,只要消费者校验出非法值,就会永久阻塞在向invalidValues发送数据的代码行,消费者goroutine无法退出,wg.Wait()永远不会返回,进而导致outputStream不会被关闭,后续写输出文件的goroutine也会一直阻塞在range outputStream,整个程序卡死。
优雅处理方案
- 最优方案:如果不需要保留非法值,直接删除
invalidValues通道,非法值跳过不发送,从根源避免阻塞问题 - 如果需要保留非法值通道做统计但不需要落地:启动一个轻量接收goroutine直接丢弃所有非法值即可:
go func() { for range invalidValues { // 空逻辑直接丢弃数据,避免发送方阻塞 } }()
修改后单输出通道完整示例代码
done := make(chan struct{}) defer close(done) inputStream := make(chan string) outputStream := make(chan string) //Producer reads a file with values and stores them in a channel go func() { scanner := bufio.NewScanner(file) for scanner.Scan() { inputStream <- strings.TrimSpace(scanner.Text()) } close(inputStream) }() //Consumers var wg sync.WaitGroup wg.Add(Workers) for i := 0; i < Workers; i++ { go func() { for value := range inputStream { if !checkValue(value) { time.Sleep(5 * time.Second) continue } select { case outputStream <- value: case <-done: return } time.Sleep(5 * time.Second) } wg.Done() }() } go func() { wg.Wait() close(outputStream) }() //Write outputStream file resultFile, err := os.Create("outputStream.txt") if err != nil { log.Fatal(err) } var wg2 sync.WaitGroup wg2.Add(1) go func() { for r := range outputStream { _, err := resultFile.WriteString(r + "\n") if err != nil { log.Fatal(err) } } resultFile.Close() wg2.Done() }() wg2.Wait()
内容的提问来源于stack exchange,提问作者DraQ
相关产品推荐
相关产品推荐

