You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多goroutine读取同一channel异常:仅各读一次即停止

问题:多Goroutine从同一Channel读取时提前停止的原因

我启动了两个worker goroutine从同一个channel读取值,但每个worker只读取一个值就停止了。原本期望它们会持续读取直到生产者关闭channel,但现在生产者还没关闭,发送第三个值时就阻塞了,第三个值也没被任何worker读取。

运行输出

new worker
new worker
waiting
sending 0
sending 1
sending 2
running func 1
sending value out 1
running func 0
sending value out 0

相关代码

package main

import (
    "fmt"
    "sync"
)

func workerPool(done <-chan bool, in <-chan int, numberOfWorkers int, fn func(int) int) chan int {
    out := make(chan int)
    var wg sync.WaitGroup

    for i := 0; i < numberOfWorkers; i++ {
        fmt.Println("new worker")
        wg.Add(1)
        // 启动worker从in读取数据,处理后发送到out
        go func() {
            defer wg.Done()
            for {
                select {
                case <-done:
                    fmt.Println("recieved done signal")
                    return
                case data, ok := <-in:
                    if !ok {
                        fmt.Println("no more items")
                        return
                    }
                    fmt.Println("running func", data)
                    value := fn(data)
                    fmt.Println("sending value out", value)
                    out <- value
                }
            }
        }()
    }

    fmt.Println("waiting")
    wg.Wait()
    fmt.Println("done waiting")
    close(out)
    return out
}

func main() {
    done := make(chan bool)
    defer close(done)

    in := make(chan int)

    go func() {
        for i := 0; i < 10; i++ {
            fmt.Println("sending", i)
            in <- i
        }
        close(in)
    }()

    out := workerPool(done, in, 2, func(i int) int {
        return i
    })

    for {
        select {
        case o, ok := <-out:
            if !ok {
                continue
            }

            fmt.Println("output", o)
        case <-done:
            return
        default:
        }
    }

}

问题原因与解决方案

核心问题:双重阻塞导致死锁

  1. Worker阻塞在结果发送:
    out是无缓冲channel,worker处理完数据执行out <- value时,需要等待接收方(main函数)读取数据才能完成发送。但workerPool启动goroutine后立刻调用wg.Wait(),导致main函数无法进入读取out的循环,worker永远阻塞在out <- value,无法继续读取in的后续值。

  2. wg.Wait()位置错误:
    workerPool函数内直接调用wg.Wait(),但此时所有worker都阻塞在发送操作上,永远无法执行wg.Done(),导致wg.Wait()无限阻塞,workerPool无法返回,main函数彻底卡死在调用workerPool的步骤。

修复方案

方案1:调整等待逻辑,异步等待worker完成

修改workerPool,不要在函数主线程调用wg.Wait(),而是启动单独goroutine等待,让workerPool立刻返回out,让main函数可以开始读取结果:

func workerPool(done <-chan bool, in <-chan int, numberOfWorkers int, fn func(int) int) chan int {
    out := make(chan int)
    var wg sync.WaitGroup

    for i := 0; i < numberOfWorkers; i++ {
        fmt.Println("new worker")
        wg.Add(1)
        go func() {
            defer wg.Done()
            for {
                select {
                case <-done:
                    fmt.Println("recieved done signal")
                    return
                case data, ok := <-in:
                    if !ok {
                        fmt.Println("no more items")
                        return
                    }
                    fmt.Println("running func", data)
                    value := fn(data)
                    fmt.Println("sending value out", value)
                    out <- value
                }
            }
        }()
    }

    // 异步等待所有worker完成后关闭out
    go func() {
        wg.Wait()
        fmt.Println("done waiting")
        close(out)
    }()

    fmt.Println("worker pool started")
    return out
}

方案2:简化main函数的结果读取逻辑

main函数里的读取循环可以直接遍历out,当out被关闭时循环会自动退出,无需额外的select判断:

func main() {
    done := make(chan bool)
    defer close(done)

    in := make(chan int)

    go func() {
        for i := 0; i < 10; i++ {
            fmt.Println("sending", i)
            in <- i
        }
        close(in)
    }()

    out := workerPool(done, in, 2, func(i int) int {
        return i
    })

    // 直接遍历out,直到channel关闭
    for o := range out {
        fmt.Println("output", o)
    }
}

内容的提问来源于stack exchange,提问作者Jonathan Kittell

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 13:20:01