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

Go协程从通道读取输入并写入另一通道时程序阻塞的问题求助

Go协程从通道读取输入并写入另一通道时程序阻塞的问题求助

你好,我来帮你分析下代码里的问题,以及对应的解决办法:

问题分析

你的代码出现阻塞和无FINISHED日志的情况,主要有两个核心原因:

  1. 无缓冲通道导致发送阻塞
    你创建的responseChan是无缓冲通道(make(chan int)),当worker协程执行responseChan <- resp.StatusCode时,必须等到主程序准备好接收这个值,worker才能继续往下执行。但你的主程序是先把所有25个url都发送到urlChan之后,才开始读取responseChan的。这就导致worker在处理完第一个url后,就会卡在发送响应码的步骤,无法继续处理后续url,也无法打印FINISHED日志。

  2. 未处理HTTP请求错误
    代码里resp, _ := http.Get(url)直接忽略了错误,如果某个url请求失败(比如网络问题、域名不存在),resp会是nil,此时访问resp.StatusCode会直接触发panic,导致该worker协程崩溃。崩溃后虽然defer wg.Done()会执行,但这个worker没有往responseChan发送值,主程序循环读取25次的要求无法满足,最终会卡在读取responseChan的步骤。

另外还有一个小问题:你把defer wg.Done()放在了for循环之后,虽然逻辑上没问题,但更稳妥的做法是把它放在函数开头,确保无论函数以什么方式退出,都能正确调用wg.Done()。

修复方案

方案1:给responseChan添加缓冲

把responseChan改成有缓冲的通道,缓冲大小至少设为最大协程数,或者直接设为url的总数,这样worker发送响应码时不会立即阻塞,可以继续处理后续请求:

responseChan := make(chan int, len(urls)) // 或者设为MaxGoroutineCount

方案2:异步接收响应码

在主程序中启动一个单独的协程来读取responseChan,这样主程序可以一边发送url到urlChan,一边接收响应,避免worker阻塞:

// 启动协程接收响应
go func() {
    for code := range responseChan {
        fmt.Println("收到响应码:", code)
    }
}()

// 填充urlChan
for _, url := range urls {
    urlChan <- url
}
close(urlChan)

wg.Wait()
close(responseChan)

方案3:处理HTTP请求错误

在worker中添加错误处理,避免panic,同时在请求失败时发送一个标识性的状态码(比如0)到通道:

func downloaderInputUrlFromChannel(id int, urlChan chan string, responseChan chan int, wg *sync.WaitGroup) {
    defer wg.Done() // 移到函数开头,确保必执行
    for url := range urlChan {
        resp, err := http.Get(url)
        if err != nil {
            fmt.Println("worker", id, "处理URL失败", url, "错误:", err)
            responseChan <- 0 // 发送错误标识
            continue
        }
        defer resp.Body.Close() // 别忘了关闭响应体,避免资源泄漏
        fmt.Println("STARTED worker ", id, " For URL ", url, "With Response Code", resp.StatusCode)
        responseChan <- resp.StatusCode
        fmt.Println("FINISHED worker ", id, " For URL ", url, "With Response Code", resp.StatusCode)
    }
}

完整修复后的代码示例

package main

import (
    "fmt"
    "net/http"
    "sync"
)

func downloaderInputUrlFromChannel(id int, urlChan chan string, responseChan chan int, wg *sync.WaitGroup) {
    defer wg.Done()
    for url := range urlChan {
        resp, err := http.Get(url)
        if err != nil {
            fmt.Println("FAILED worker", id, "For URL", url, "Error:", err)
            responseChan <- 0
            continue
        }
        defer resp.Body.Close()
        fmt.Println("STARTED worker ", id, " For URL ", url, "With Response Code", resp.StatusCode)
        responseChan <- resp.StatusCode
        fmt.Println("FINISHED worker ", id, " For URL ", url, "With Response Code", resp.StatusCode)
    }
}

func main() {
    urls := []string{
        "https://httpbin.org?q=1,l=2,p=1",
        "https://httpbin.org/",
        "https://httpbin.org?q=1,r=2,p=1",
        "https://httpbin.org?q=1",
        "https://httpbin.org?a=1",
        // 补充剩余20个url
    }
    MaxGoroutineCount := 5
    var wg sync.WaitGroup

    urlChan := make(chan string, MaxGoroutineCount) // 给urlChan加缓冲,提高发送效率
    responseChan := make(chan int, MaxGoroutineCount)

    // 启动Workers
    for i := 0; i < MaxGoroutineCount; i++ {
        wg.Add(1)
        go downloaderInputUrlFromChannel(i, urlChan, responseChan, &wg)
    }

    // 异步接收响应
    go func() {
        for i := 0; i < len(urls); i++ {
            fmt.Println("Received response code:", <-responseChan)
        }
        close(responseChan)
    }()

    // 填充urlChan
    for _, url := range urls {
        urlChan <- url
    }
    close(urlChan)

    wg.Wait()
}

备注:内容来源于stack exchange,提问作者Ishan Bhatt

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.20 10:34:31