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

OpenAI流式API请求超时处理及http2响应体关闭错误解决

流式OpenAI API请求超时处理问题

问题背景

我正在开发一款使用OpenAI流式聊天补全API的应用,流式功能原本运行正常,但希望在API未在指定时间响应时处理超时逻辑。添加超时检查后,出现**"http2: response body closed"**错误,无法正常接收响应。

无超时检查的正常代码

func (g *gptAdaptorClient) CreateChatCompletionStream(messages []map[string]string, model string, maxTokens int, temperature float64) (<-chan string, error) {
    url := "https://api.openai.com/v1/chat/completions"

    // Payload to send to the API
    payload := map[string]interface{}{
        "model":       model,
        "messages":    messages,
        "max_tokens":  maxTokens,
        "temperature": temperature,
        "stream":      true, // Request to return as a stream
    }

    requestBody, err := json.Marshal(payload)
    if err != nil {
        return nil, fmt.Errorf("failed to marshal payload: %w", err)
    }

    // Create HTTP request
    req, err := http.NewRequest("POST", url, bytes.NewBuffer(requestBody))
    if err != nil {
        return nil, fmt.Errorf("failed to create request: %w", err)
    }

    g.addCommonHeaders(req) // Add necessary headers

    // Send request and get response
    resp, err := g.Client.Do(req)
    if err != nil {
        return nil, fmt.Errorf("failed to send request: %w", err)
    }

    // Check HTTP status code
    if resp.StatusCode != http.StatusOK {
        body, _ := io.ReadAll(resp.Body)
        return nil, fmt.Errorf("API error: %s", string(body))
    }

    // Create channel to pass stream data
    dataChannel := make(chan string)

    // Process stream in a goroutine
    go func() {
        defer close(dataChannel)
        defer resp.Body.Close()

        scanner := bufio.NewScanner(resp.Body)
        for scanner.Scan() {
            line := scanner.Text()

            // Check if the line doesn't contain data
            if len(line) < 6 || line[:6] != "data: " {
                continue
            }

            // Extract the JSON content after "data: "
            chunk := line[6:]
            if chunk == "[DONE]" {
                break
            }

            // Send chunk to the channel
            dataChannel <- chunk
        }

        // Check for scanner errors (if any)
        if err := scanner.Err(); err != nil {
            fmt.Printf("Error reading streaming response: %v\n", err)
        }
    }()

    return dataChannel, nil
}

添加超时后报错的代码

func (g *gptAdaptorClient) CreateChatCompletionStream(messages []map[string]string, model string, maxTokens int, temperature float64) (<-chan string, error) {
    url := "https://api.openai.com/v1/chat/completions"

    // Payload to send to the API
    payload := map[string]interface{}{
        "model":       model,
        "messages":    messages,
        "max_tokens":  maxTokens,
        "temperature": temperature,
        "stream":      true, // Request to return as a stream
    }

    requestBody, err := json.Marshal(payload)
    if err != nil {
        return nil, fmt.Errorf("failed to marshal payload: %w", err)
    }

    // Create HTTP request
    req, err := http.NewRequest("POST", url, bytes.NewBuffer(requestBody))
    if err != nil {
        return nil, fmt.Errorf("failed to create request: %w", err)
    }

    g.addCommonHeaders(req) // Add necessary headers

    // Set up the context with timeout for HTTP request (10 seconds)
    reqCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()

    req = req.WithContext(reqCtx) // Send the request with the timeout context

    // Send request and get response
    resp, err := g.Client.Do(req)
    if err != nil {
        // Check if the error is a timeout and switch to Perplexity API
        if reqCtx.Err() == context.DeadlineExceeded {
            fmt.Println("OpenAI API call timed out, switching to Perplexity API...")
            return g.callPerplexityApi(messages, model, maxTokens, temperature)
        }
        return nil, fmt.Errorf("failed to send request: %w", err)
    }
    defer resp.Body.Close()

    // Check HTTP status code
    if resp.StatusCode != http.StatusOK {
        body, _ := io.ReadAll(resp.Body)
        return nil, fmt.Errorf("OpenAI API error: %s", string(body))
    }

    // Create channel to pass stream data
    dataChannel := make(chan string)

    // Process stream in a goroutine
    go func() {
        defer close(dataChannel)

        // Use scanner to read lines from the response body
        scanner := bufio.NewScanner(resp.Body)
        for scanner.Scan() {
            line := scanner.Text()

            // Check if the line doesn't contain data
            if len(line) < 6 || line[:6] != "data: " {
                continue
            }

            // Extract the JSON content after "data: "
            chunk := line[6:]
            if chunk == "[DONE]" {
                break
            }

            // Send chunk to the channel
            dataChannel <- chunk
        }

        // Check for scanner errors (if any)
        if err := scanner.Err(); err != nil {
            fmt.Printf("Error reading streaming response from OpenAI API: %v\n", err)
        }
    }()

    return dataChannel, nil
}

问题详情

  • 无超时版本:流式功能可正常运行
  • 超时引发的问题:为HTTP请求添加10秒上下文超时后,触发**"http2: response body closed"**错误,无法获取OpenAI API响应
  • 预期目标:实现OpenAI API 10秒未响应时,自动切换至Perplexity等备用API

技术疑问

  1. 为何为HTTP请求添加上下文超时会在流式场景下引发**"http2: response body closed"**错误?
  2. 如何在添加请求超时的同时避免该问题,确保请求成功后流式功能正常运行?
  3. 在Go中处理OpenAI这类流式API的HTTP请求超时,是否有更优方案可避免提前关闭响应体?

解答

1. 错误原因分析

你添加的context.WithTimeout会在10秒后触发上下文取消,而defer cancel()会在函数退出时立即取消上下文。问题在于:当请求成功拿到响应后,主函数退出时执行defer cancel(),此时上下文被取消,而HTTP/2客户端会监听上下文状态,一旦上下文取消就会主动关闭响应体连接,导致后续goroutine读取流时出现http2: response body closed错误。

简单来说:主函数返回后,defer cancel()执行,上下文失效,HTTP客户端直接断开了流式响应的连接。

2. 解决方案:分离请求超时与流式读取的上下文

我们需要让请求阶段的超时只作用于获取响应前,拿到响应后就取消超时上下文的绑定,改用独立的上下文管理流式读取过程。

修改后的代码示例:

func (g *gptAdaptorClient) CreateChatCompletionStream(messages []map[string]string, model string, maxTokens int, temperature float64) (<-chan string, error) {
    url := "https://api.openai.com/v1/chat/completions"

    payload := map[string]interface{}{
        "model":       model,
        "messages":    messages,
        "max_tokens":  maxTokens,
        "temperature": temperature,
        "stream":      true,
    }

    requestBody, err := json.Marshal(payload)
    if err != nil {
        return nil, fmt.Errorf("failed to marshal payload: %w", err)
    }

    req, err := http.NewRequest("POST", url, bytes.NewBuffer(requestBody))
    if err != nil {
        return nil, fmt.Errorf("failed to create request: %w", err)
    }

    g.addCommonHeaders(req)

    // 仅为请求阶段(获取响应)设置超时上下文
    reqCtx, cancelReq := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancelReq() // 仅在请求阶段结束后取消,拿到响应后这个cancel不影响流式读取

    req = req.WithContext(reqCtx)

    resp, err := g.Client.Do(req)
    if err != nil {
        if reqCtx.Err() == context.DeadlineExceeded {
            fmt.Println("OpenAI API call timed out, switching to Perplexity API...")
            return g.callPerplexityApi(messages, model, maxTokens, temperature)
        }
        return nil, fmt.Errorf("failed to send request: %w", err)
    }

    if resp.StatusCode != http.StatusOK {
        body, _ := io.ReadAll(resp.Body)
        resp.Body.Close()
        return nil, fmt.Errorf("OpenAI API error: %s", string(body))
    }

    dataChannel := make(chan string)

    // 流式读取使用独立的上下文,不受请求超时上下文影响
    go func() {
        defer close(dataChannel)
        defer resp.Body.Close()

        scanner := bufio.NewScanner(resp.Body)
        for scanner.Scan() {
            line := scanner.Text()

            if len(line) < 6 || line[:6] != "data: " {
                continue
            }

            chunk := line[6:]
            if chunk == "[DONE]" {
                break
            }

            dataChannel <- chunk
        }

        if err := scanner.Err(); err != nil {
            fmt.Printf("Error reading streaming response from OpenAI API: %v\n", err)
        }
    }()

    return dataChannel, nil
}

关键修改点:

  • defer cancelReq()在主函数退出时执行,但此时已经拿到响应,HTTP客户端不会因为这个上下文取消而关闭响应体(因为请求已经完成,响应阶段的连接不受请求上下文的影响)
  • 流式读取的goroutine独立运行,响应体的关闭由goroutine内部的defer resp.Body.Close()负责,直到流读取完成或出错

3. 更优方案:精细化超时控制

对于流式API,建议将超时分为两个阶段:

  • 请求建立超时:控制从发送请求到收到响应头的时间(比如10秒),这就是你需要的超时逻辑,用来触发备用API切换
  • 流读取超时:控制流式响应中两个数据块之间的间隔时间,避免连接挂起但无数据传输的情况

实现流读取超时的示例:

go func() {
    defer close(dataChannel)
    defer resp.Body.Close()

    reader := &bufio.Reader{Reader: resp.Body}
    scanner := bufio.NewScanner(reader)

    for {
        // 为每次扫描设置超时
        err := reader.SetReadDeadline(time.Now().Add(30 * time.Second))
        if err != nil {
            fmt.Printf("Failed to set read deadline: %v\n", err)
            return
        }

        if !scanner.Scan() {
            break
        }

        line := scanner.Text()
        if len(line) < 6 || line[:6] != "data: " {
            continue
        }

        chunk := line[6:]
        if chunk == "[DONE]" {
            break
        }

        dataChannel <- chunk
    }

    if err := scanner.Err(); err != nil {
        fmt.Printf("Error reading streaming response from OpenAI API: %v\n", err)
    }
}()

这样既保证了请求阶段的超时切换逻辑,又能防止流式连接长时间无数据传输的情况,同时避免了提前关闭响应体的问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 15:15:56