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
技术疑问
- 为何为HTTP请求添加上下文超时会在流式场景下引发**"http2: response body closed"**错误?
- 如何在添加请求超时的同时避免该问题,确保请求成功后流式功能正常运行?
- 在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
相关产品推荐
相关产品推荐

