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

使用WaitGroups与Buffered Channels的Go代码死锁原因排查

WaitGroups、Buffered Channels与死锁

这段Go代码出现死锁,需求是通过inputChan发送数据,再从outputChan读取数据。

原始代码:

package main

import (
    "fmt"
    "sync"
)

func listStuff(wg *sync.WaitGroup, workerID int, inputChan chan int, outputChan chan int) {
    defer wg.Done()

    for i := range inputChan {
        fmt.Println("sending ", i)
        outputChan <- i
    }
}

func List(workers int) ([]int, error) {
    _output := make([]int, 0)

    inputChan := make(chan int, 1000)
    outputChan := make(chan int, 1000)

    var wg sync.WaitGroup
    wg.Add(workers)

    fmt.Printf("+++ Spinning up %v workers\n", workers)
    for i := 0; i < workers; i++ {
        go listStuff(&wg, i, inputChan, outputChan)
    }

    for i := 0; i < 3000; i++ {
        inputChan <- i
    }

    done := make(chan struct{})
    go func() {
        close(done)
        close(inputChan)
        close(outputChan)
        wg.Wait()
    }()

    for o := range outputChan {
        fmt.Println("reading from channel...")
        _output = append(_output, o)
    }

    <-done
    fmt.Printf("+++ output len: %v\n", len(_output))
    return _output, nil
}

func main() {
    List(5)
}

死锁原因

  1. 双向阻塞:主goroutine向inputChan发送3000个数据,当inputChan缓冲(1000)填满后,需要等待worker取走数据才能继续发送;而worker从inputChan取数据后向outputChan发送,当outputChan缓冲(1000)填满后,worker阻塞,无法继续读取inputChan,导致主goroutine也阻塞,形成死锁。
  2. 通道关闭顺序错误:匿名goroutine提前关闭outputChan,不仅会导致worker向已关闭通道发送数据触发panic,还无法保证worker完成所有数据转发。

修复代码

package main

import (
    "fmt"
    "sync"
)

func listStuff(wg *sync.WaitGroup, workerID int, inputChan chan int, outputChan chan int) {
    defer wg.Done()

    for i := range inputChan {
        fmt.Println("worker", workerID, "sending", i)
        outputChan <- i
    }
}

func List(workers int) ([]int, error) {
    _output := make([]int, 0)

    inputChan := make(chan int, 1000)
    outputChan := make(chan int, 1000)

    var wg sync.WaitGroup
    wg.Add(workers)

    fmt.Printf("+++ Spinning up %v workers\n", workers)
    for i := 0; i < workers; i++ {
        go listStuff(&wg, i, inputChan, outputChan)
    }

    // 单独goroutine发送数据,避免主goroutine阻塞在发送操作上
    go func() {
        for i := 0; i < 3000; i++ {
            inputChan <- i
        }
        // 发送完所有数据后关闭inputChan,告知worker没有更多数据
        close(inputChan)
    }()

    // 等待所有worker完成后关闭outputChan
    go func() {
        wg.Wait()
        close(outputChan)
    }()

    // 从outputChan读取所有数据,直到通道关闭
    for o := range outputChan {
        fmt.Println("reading from channel...", o)
        _output = append(_output, o)
    }

    fmt.Printf("+++ output len: %v\n", len(_output))
    return _output, nil
}

func main() {
    List(5)
}

修复说明

  1. 异步发送数据:将向inputChan发送数据的逻辑放到单独goroutine中,避免主goroutine阻塞在发送操作,保证主goroutine可以持续从outputChan读取数据,缓解通道缓冲满的问题。
  2. 正确的通道关闭顺序:
    • 数据发送完成后立即关闭inputChan,让worker的for range循环在读取完所有数据后退出。
    • 等待所有worker完成(wg.Wait())后再关闭outputChan,确保所有数据都已转发到outputChan,主goroutine可以读取到全部数据。
  3. 消除死锁:主goroutine持续读取outputChan,避免worker因outputChan缓冲满而阻塞,worker就能持续从inputChan取数据,保证数据流动不中断。

内容的提问来源于stack exchange,提问作者salman b

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 03:05:19