Go语言中使用Channel并发接收API响应并写入SQL的问题
Go并发API数据入库的通道问题解决
我用Go实现一套数据处理流程:从外部API获取JSON数据,处理后写入SQL数据库。尝试并发发起API请求,拿到响应后想用另一个goroutine调用load()函数插入数据库,但代码运行时有时能看到load()里的log.Printf()输出,有时看不到,推测是通道关闭或通信逻辑有问题。以下是我的代码:
package main import ( "encoding/json" "io/ioutil" "log" "net/http" "time" ) type Request struct { url string } type Response struct { status int args Args `json:"args"` headers Headers `json:"headers"` origin string `json:"origin"` url string `json:"url"` } type Args struct { } type Headers struct { accept string `json:"Accept"` } func main() { start := time.Now() numRequests := 5 responses := make(chan Response, 5) defer close(responses) for i := 0; i < numRequests; i++ { req := Request{url: "https://httpbin.org/get"} go func(req *Request) { resp, err := extract(req) if err != nil { log.Fatal("Error extracting data from API") return } // Send response to channel responses <- resp }(&req) // Perform go routine to load data go load(responses) } log.Println("Execution time: ", time.Since(start)) } func extract(req *Request) (r Response, err error) { var resp Response request, err := http.NewRequest("GET", req.url, nil) if err != nil { return resp, err } request.Header = http.Header{ "accept": {"application/json"}, } response, err := http.DefaultClient.Do(request) defer response.Body.Close() if err != nil { log.Fatal("Error") return resp, err } // Read response data body, err := ioutil.ReadAll(response.Body) if err != nil { log.Fatal("Error") return resp, err } json.Unmarshal(body, &resp) resp.status = response.StatusCode return resp, nil } type Record struct { origin string url string } func load(ch chan Response) { // Read response from channel resp := <-ch // Process the response data records := process(resp) log.Printf("%+v\n", records) // Load data to db stuff here } func process(resp Response) (record Record) { // Process the response struct as needed to get a record of data to insert to DB return record }
问题分析
- main函数提前退出:main函数启动所有goroutine后立刻打印执行时间并退出,此时Go程序会终止所有未完成的goroutine,导致部分
load()还没执行就被终止,日志时有时无。 - 通道与goroutine匹配逻辑错误:循环中每发起一个请求就启动一个
load()goroutine,共5个,但通道最多存5个响应。若某个load先启动但通道无数据,main退出时它会被直接终止;同时defer close(responses)在main退出时执行,此时可能还有extractgoroutine尝试往通道发数据,会触发panic。 process函数未有效赋值:当前process返回零值Record,即使执行也看不到有效数据,易误导判断。
解决方案
- 等待所有goroutine完成:用
sync.WaitGroup跟踪请求和入库goroutine的执行状态,确保main在所有任务完成后再退出。 - 改用worker池处理入库:启动固定数量的worker goroutine,持续从通道读取数据并处理,避免每个请求对应一个
load的资源浪费,同时保证所有响应都被处理。 - 正确关闭通道:在所有请求goroutine完成后再关闭通道,避免发送数据到已关闭通道的panic。
- 修复
process函数:正确从Response提取数据到Record,确保日志输出有效内容。
修正后的完整代码
package main import ( "encoding/json" "io/ioutil" "log" "net/http" "sync" "time" ) type Request struct { url string } type Response struct { status int args Args `json:"args"` headers Headers `json:"headers"` origin string `json:"origin"` url string `json:"url"` } type Args struct{} type Headers struct { accept string `json:"Accept"` } func main() { start := time.Now() numRequests := 5 numWorkers := 5 // 可根据实际情况调整,比如与CPU核心数一致 responses := make(chan Response, numRequests) var wgExtract sync.WaitGroup var wgLoad sync.WaitGroup // 启动worker池处理入库 wgLoad.Add(numWorkers) for i := 0; i < numWorkers; i++ { go func() { defer wgLoad.Done() for resp := range responses { load(resp) } }() } // 并发发起API请求 wgExtract.Add(numRequests) for i := 0; i < numRequests; i++ { req := Request{url: "https://httpbin.org/get"} go func(req Request) { defer wgExtract.Done() resp, err := extract(&req) if err != nil { log.Printf("Error extracting data from API: %v", err) return } responses <- resp }(req) // 直接传值避免闭包引用同一变量的问题 } // 等待所有请求完成后关闭通道 wgExtract.Wait() close(responses) // 等待所有worker处理完成 wgLoad.Wait() log.Println("Execution time: ", time.Since(start)) } func extract(req *Request) (r Response, err error) { var resp Response request, err := http.NewRequest("GET", req.url, nil) if err != nil { return resp, err } request.Header.Set("accept", "application/json") response, err := http.DefaultClient.Do(request) if err != nil { return resp, err } defer response.Body.Close() // 移到err判断后,避免response为nil时调用Close() body, err := ioutil.ReadAll(response.Body) if err != nil { return resp, err } if err := json.Unmarshal(body, &resp); err != nil { return resp, err } resp.status = response.StatusCode return resp, nil } type Record struct { origin string url string } func load(resp Response) { records := process(resp) log.Printf("%+v\n", records) // 这里添加数据库插入逻辑 } func process(resp Response) Record { // 正确提取数据 return Record{ origin: resp.origin, url: resp.url, } }
关键修正点
- 用
sync.WaitGroup分别跟踪请求和worker goroutine,确保所有任务完成后main才退出。 - 改用worker池模式,固定数量的worker从通道循环读取数据,直到通道关闭。
- 将
defer response.Body.Close()移到err != nil判断之后,避免response为nil时触发panic。 - 修复闭包引用问题:goroutine传值
req而非指针,避免循环中变量覆盖的问题。 process函数现在正确返回提取后的Record,日志能输出有效内容。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

