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

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退出时执行,此时可能还有extract goroutine尝试往通道发数据,会触发panic。
  • process函数未有效赋值:当前process返回零值Record,即使执行也看不到有效数据,易误导判断。

解决方案

  1. 等待所有goroutine完成:用sync.WaitGroup跟踪请求和入库goroutine的执行状态,确保main在所有任务完成后再退出。
  2. 改用worker池处理入库:启动固定数量的worker goroutine,持续从通道读取数据并处理,避免每个请求对应一个load的资源浪费,同时保证所有响应都被处理。
  3. 正确关闭通道:在所有请求goroutine完成后再关闭通道,避免发送数据到已关闭通道的panic。
  4. 修复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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:42:04