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

Go通道通信程序执行完毕后挂起,请求排查原因

Go并发程序挂起问题排查与修复

问题根源

你的程序运行结束后挂起的核心原因是主线程阻塞在respChan的for range循环中,无法执行后续的wg.Wait()和close(respChan)操作:

  • 启动GET协程后,主线程立刻进入for ele := range respChan循环,该循环会持续等待通道接收数据,直到通道被关闭。
  • 而wg.Wait()(等待所有GET协程完成)和close(respChan)的代码位于循环之后,永远无法被执行。
  • 当所有GET协程都完成数据发送后,respChan不再有新数据,但通道未关闭,for range会一直阻塞,导致程序挂起。

额外优化点

原代码中POST请求的协程没有被等待,主线程可能在POST请求完成前就退出,导致部分POST请求未执行完成;同时HTTP响应体未关闭,存在资源泄漏风险。

修复方案

  1. 将respChan的遍历逻辑放到独立协程中,让主线程可以继续执行wg.Wait(),等待所有GET协程完成后关闭通道。
  2. 新增postWg等待组,确保所有POST协程完成后主线程再退出。
  3. 在HTTP请求处理函数中添加响应体关闭逻辑,避免资源泄漏。

修正后的完整代码

package main

import (
    "bytes"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "sync"
    "time"
)

type HttpBinGetRequest struct {
    url string
}

type HttpBinGetResponse struct {
    Uuid       string `json:"uuid"`
    StatusCode int
}

type HttpBinPostRequest struct {
    url  string
    uuid string // Item to post to API
}

type HttpBinPostResponse struct {
    Data       string `json:"data"`
    StatusCode int
}

func main() {
    // Prepare GET requests for n requests
    var requests []*HttpBinGetRequest
    for i := 0; i < 10; i++ {
        uri := "https://httpbin.org/uuid"
        request := &HttpBinGetRequest{
            url: uri,
        }
        requests = append(requests, request)
    }

    // Create semaphore and rate limit for the GET endpoint
    getSemaphore := make(chan struct{}, 10)
    getRate := make(chan struct{}, 10)
    defer close(getRate)
    defer close(getSemaphore)
    for i := 0; i < cap(getRate); i++ {
        getRate <- struct{}{}
    }

    go func() {
        ticker := time.NewTicker(100 * time.Millisecond)
        defer ticker.Stop()
        for range ticker.C {
            _, ok := <-getRate
            if !ok {
                return
            }
        }
    }()

    // Send our GET requests to obtain a random UUID
    respChan := make(chan HttpBinGetResponse)
    var getWg sync.WaitGroup
    for _, request := range requests {
        getWg.Add(1)
        go func(r *HttpBinGetRequest) {
            defer getWg.Done()

            getRate <- struct{}{}
            resp, _ := get(r, getSemaphore)

            fmt.Printf("GET Response: %+v\n", resp)
            respChan <- *resp
        }(request)
    }

    // Set up for POST requests 10/s
    postSemaphore := make(chan struct{}, 10)
    postRate := make(chan struct{}, 10)
    defer close(postRate)
    defer close(postSemaphore)
    for i := 0; i < cap(postRate); i++ {
        postRate <- struct{}{}
    }

    go func() {
        ticker := time.NewTicker(100 * time.Millisecond)
        defer ticker.Stop()
        for range ticker.C {
            _, ok := <-postRate
            if !ok {
                return
            }
        }
    }()

    // 新增POST等待组,确保所有POST请求完成
    var postWg sync.WaitGroup
    // 将respChan遍历放到独立协程,避免阻塞主线程
    go func() {
        for ele := range respChan {
            postWg.Add(1)
            postReq := &HttpBinPostRequest{
                url:  "https://httpbin.org/post",
                uuid: ele.Uuid,
            }
            go func(r *HttpBinPostRequest) {
                defer postWg.Done()
                postRate <- struct{}{}
                postResp, err := post(r, postSemaphore)
                if err != nil {
                    fmt.Println("POST Error:", err)
                }
                fmt.Printf("POST Response: %+v\n", postResp)
            }(postReq)
        }
        postWg.Wait()
    }()

    // 等待所有GET协程完成,然后关闭respChan
    getWg.Wait()
    close(respChan)
}

func get(hbgr *HttpBinGetRequest, sem chan struct{}) (*HttpBinGetResponse, error) {
    sem <- struct{}{}
    defer func() { <-sem }()

    httpResp := &HttpBinGetResponse{}
    client := &http.Client{}
    req, err := http.NewRequest("GET", hbgr.url, nil)
    if err != nil {
        fmt.Println("GET Request Error:", err)
        return httpResp, err
    }

    req.Header.Set("accept", "application/json")

    resp, err := client.Do(req)
    if err != nil {
        fmt.Println("GET Do Error:", err)
        return httpResp, err
    }
    defer resp.Body.Close() // 关闭响应体,避免资源泄漏

    body, err := io.ReadAll(resp.Body)
    if err != nil {
        fmt.Println("GET Read Body Error:", err)
        return httpResp, err
    }
    if err := json.Unmarshal(body, &httpResp); err != nil {
        fmt.Println("GET Unmarshal Error:", err)
    }
    httpResp.StatusCode = resp.StatusCode
    return httpResp, nil
}

func post(hbr *HttpBinPostRequest, sem chan struct{}) (*HttpBinPostResponse, error) {
    sem <- struct{}{}
    defer func() { <-sem }()

    httpResp := &HttpBinPostResponse{}
    client := &http.Client{}
    req, err := http.NewRequest("POST", hbr.url, bytes.NewBuffer([]byte(hbr.uuid)))
    if err != nil {
        fmt.Println("POST Request Error:", err)
        return httpResp, err
    }

    req.Header.Set("accept", "application/json")

    resp, err := client.Do(req)
    if err != nil {
        fmt.Println("POST Do Error:", err)
        return httpResp, err
    }
    defer resp.Body.Close() // 关闭响应体,避免资源泄漏

    body, err := io.ReadAll(resp.Body)
    if err != nil {
        fmt.Println("POST Read Body Error:", err)
        return httpResp, err
    }
    if err := json.Unmarshal(body, &httpResp); err != nil {
        fmt.Println("POST Unmarshal Error:", err)
    }
    httpResp.StatusCode = resp.StatusCode
    return httpResp, nil
}

关键修改说明

  1. 将respChan的for range循环移到独立协程,让主线程可以执行getWg.Wait()等待所有GET请求完成,之后关闭respChan。
  2. 新增postWg等待组,确保所有POST协程执行完毕后程序才会退出。
  3. 在get和post函数中添加defer resp.Body.Close(),修复HTTP响应体资源泄漏问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 13:58:12