
本文详解 go 客户端如何稳定接收长连接 http 流(如金融行情推送),重点解决因服务端无数据、市场休市或默认超时导致的连接中断问题,并提供可复用的带心跳保活、错误重连与上下文控制的流式请求实现。
本文详解 go 客户端如何稳定接收长连接 http 流(如金融行情推送),重点解决因服务端无数据、市场休市或默认超时导致的连接中断问题,并提供可复用的带心跳保活、错误重连与上下文控制的流式请求实现。
Go 中处理 HTTP 流式响应(如 TradeKing 的实时行情流 /v1/market/quotes.json)时,常见误区是将其当作普通短连接请求处理——调用 http.Client.Get() 后立即 ioutil.ReadAll(),这会导致整个响应体被一次性读取并关闭连接,完全违背“流”的设计意图。真正的流式接口会持续发送 chunked JSON 消息(如 { "status": "connected" }、{ "quote": { ... } }),需逐块解析、长期保持连接。
关键问题在于:
- 默认
http.Client使用net/http.DefaultTransport,其ResponseHeaderTimeout和IdleConnTimeout均为 30 秒,服务端若在休市期(如 MLK Day)不发送数据,连接将被客户端或中间代理(如 Cloudflare)静默断开; - 原代码未检查
resp.StatusCode,也未处理流式响应的分块读取逻辑,ioutil.ReadAll()会阻塞直至 EOF 或超时,而流式响应永不结束。
✅ 正确做法:使用 bufio.Scanner 按行/按块流式读取,并显式配置超时与重连:
func streamQuotes() error {
// 自定义 HTTP 客户端:禁用默认超时,启用长连接
client := &http.Client{
Transport: &http.Transport{
IdleConnTimeout: 5 * time.Minute,
ResponseHeaderTimeout: 30 * time.Second,
TLSHandshakeTimeout: 10 * time.Second,
},
}
config := oauth1.NewConfig("HIDDEN", "HIDDEN")
token := oauth1.NewToken("HIDDEN", "HIDDEN")
oauthClient := config.Client(oauth1.NoContext, token)
// 注意:此处应复用 client 而非 oauthClient 的 transport(除非它已定制)
// 实际中建议将 oauth 注入自定义 client 的 RoundTripper
path := "https://stream.tradeking.com/v1/market/quotes.json?symbols=aapl"
// 使用 context 控制整体生命周期(如 10 分钟后强制退出)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
defer cancel()
req, err := http.NewRequestWithContext(ctx, "GET", path, nil)
if err != nil {
return fmt.Errorf("failed to create request: %w", err)
}
// 手动添加 OAuth 头(更可控)
authHeader := config.AuthorizationHeader(token, "GET", path, url.Values{})
req.Header.Set("Authorization", authHeader)
resp, err := client.Do(req)
if err != nil {
return fmt.Errorf("request failed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("unexpected status code: %d", resp.StatusCode)
}
// 按行扫描流式 JSON(TradeKing 使用换行分隔的 JSON 对象)
scanner := bufio.NewScanner(resp.Body)
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
// 解析单条 JSON(如 {"status":"connected"} 或 {"quote":{...}})
var msg map[string]interface{}
if err := json.Unmarshal([]byte(line), &msg); err != nil {
log.Printf("warn: invalid JSON line: %s | error: %v", line, err)
continue
}
// 示例:识别连接状态
if status, ok := msg["status"].(string); ok && status == "connected" {
log.Println("✅ Stream connected successfully")
}
// 处理 quote 数据...
if quote, ok := msg["quote"]; ok {
log.Printf("? Quote received: %+v", quote)
}
}
if err := scanner.Err(); err != nil {
// 网络中断、超时等错误
log.Printf("⚠️ Stream ended with error: %v", err)
return err
}
return nil
}
⚠️ 重要注意事项:
-
不要依赖
oauth1.Client的默认 transport:它未针对流式场景优化,应提取其认证逻辑,注入到自定义http.Client中; -
务必校验
StatusCode:休市时 API 可能返回204 No Content或403 Forbidden,而非空响应; -
添加重试机制:生产环境需封装
streamQuotes()为带指数退避的循环(如for i := 0; i ); -
监控连接健康度:可启动 goroutine 定期发送
PING(若服务端支持)或记录最后接收时间,超时则主动重连; -
资源清理:始终用
defer resp.Body.Close(),并在context取消时确保 goroutine 退出。
通过以上改造,你的 Go 服务即可稳健维持与行情流的长连接,从容应对市场开闭市、网络抖动等真实场景。










