Go协程从通道读取输入并写入另一通道时程序阻塞的问题求助
你好,我来帮你分析下代码里的问题,以及对应的解决办法:
问题分析
你的代码出现阻塞和无FINISHED日志的情况,主要有两个核心原因:
无缓冲通道导致发送阻塞
你创建的responseChan是无缓冲通道(make(chan int)),当worker协程执行responseChan <- resp.StatusCode时,必须等到主程序准备好接收这个值,worker才能继续往下执行。但你的主程序是先把所有25个url都发送到urlChan之后,才开始读取responseChan的。这就导致worker在处理完第一个url后,就会卡在发送响应码的步骤,无法继续处理后续url,也无法打印FINISHED日志。未处理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

