推理服务的API与普通REST接口有一个明显区别:单次请求可能要等待数秒甚至数十秒,返回的token还可能以流式方式持续推送。直接用Go标准库的http.Client写一个简单封装,功能上没有问题,但一旦压到每秒几百个并发请求,就会暴露出连接复用率低、失败重试无节制、流式读取阻塞goroutine等一系列问题。一个能扛住高并发的推理API客户端SDK,核心并不在于请求逻辑本身,而在于传输层、并发控制、错误恢复和流式消费这四个层面的配合。

传输层连接池与并发上限设计
Go的net/http默认Transport在连接池参数上偏向保守,尤其是每个主机的空闲连接数MaxIdleConnsPerHost默认只有2。对于高并发场景,这个值会导致大量连接被频繁创建和销毁,TLS握手开销被无限放大。推理API通常走HTTPS,握手成本更高,所以第一步就要把连接池参数调大。与此同时,还要限制最大连接数,避免客户端把推理服务的文件描述符打满。实际调优时,MaxIdleConnsPerHost可以设置为单机预期并发的50%到100%,MaxConnsPerHost则设置一个硬上限,超过这个值的请求会等待连接释放或者直接快速失败。
除了连接池本身,DialContext的超时、TLS握手超时和响应头超时也必须显式设置。默认的http.Client没有超时,单个请求可能无限期挂起。推理API在负载升高时响应时间会明显变长,但客户端必须有一个明确的等待边界。下面是一段初始化专用Transport和Client的代码,参数可以根据压测结果调整。
package inference
import (
"net"
"net/http"
"time"
)
func newHTTPClient() *http.Client {
transport := &http.Transport{
Proxy: http.ProxyFromEnvironment,
DialContext: (&net.Dialer{
Timeout: 5 * time.Second,
KeepAlive: 30 * time.Second,
}).DialContext,
MaxIdleConns: 200,
MaxIdleConnsPerHost: 100,
MaxConnsPerHost: 150,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 5 * time.Second,
ExpectContinueTimeout: 1 * time.Second,
ForceAttemptHTTP2: true,
}
return &http.Client{
Transport: transport,
Timeout: 0, // 请求级超时由context控制,不要在这里设置全局超时
}
}
需要注意的是,http.Client的Timeout字段是整次请求的总超时,包括连接建立、发送请求和读取响应。对于流式推理,读取阶段可能持续几十秒,直接给Client设置Timeout会导致流正常返回时被误杀。更合理的做法是只在context中控制超时,并针对不同阶段设置不同值。
并发上限不能只依赖MaxConnsPerHost,因为当连接数达到上限时,http.Client会阻塞等待连接,这个等待可能没有超时。更可控的方式是在SDK内部增加一个带缓冲的信号量channel,把在途请求数限制在一个经过压测验证的安全水平。例如推理服务单实例峰值支持200并发,那么SDK端可以配置为150,留出余量。超过限制时直接返回错误,而不是无限排队。这样做既保护了后端,也能让调用方快速感知压力。
超时控制与指数退避重试
推理API的超时不能只设一个固定值。短文本生成的请求可能在1秒内完成,长文本或复杂推理可能需要30秒。如果统一使用30秒超时,那么大量失败请求会占用连接资源和goroutine,拖垮整个客户端。更好的做法是让调用方通过context传递每次调用的期望超时,SDK内部只负责把context超时映射到http.Request上。同时,连接建立、响应头到达、流式数据间隔这三个阶段可以分别设置更短的时限。
重试策略是推理API封装中最容易出问题的部分。网络抖动和5xx响应可以重试,但推理服务可能因为某个请求占用了过多GPU内存导致超时,重试同样的请求大概率还是失败。因此,重试必须是非幂等安全的。如果API支持幂等键,SDK可以自动为每次请求生成并复用同一个键,保证重试不会导致重复推理。下面这段代码展示了基础重试逻辑,只对网络错误、429和5xx进行有限次重试,并加入指数退避和随机抖动。
package inference
import (
"context"
"errors"
"math"
"math/rand"
"net/http"
"time"
)
func (c *Client) doWithRetry(ctx context.Context, req *http.Request) (*http.Response, error) {
maxRetries := 3
baseDelay := 200 * time.Millisecond
for attempt := 0; attempt <= maxRetries; attempt++ {
resp, err := c.httpClient.Do(req.Clone(ctx))
if err == nil && resp.StatusCode < 500 && resp.StatusCode != http.StatusTooManyRequests {
return resp, nil
}
if resp != nil {
// 4xx通常表示请求本身有问题,重试没有意义
if resp.StatusCode >= 400 && resp.StatusCode < 500 && resp.StatusCode != http.StatusTooManyRequests {
return resp, nil
}
resp.Body.Close()
}
if attempt == maxRetries {
if err != nil {
return nil, err
}
return resp, nil
}
backoff := float64(baseDelay) * math.Pow(2, float64(attempt))
jitter := time.Duration(rand.Float64() * float64(baseDelay))
select {
case <-time.After(time.Duration(backoff) + jitter):
case <-ctx.Done():
return nil, ctx.Err()
}
}
return nil, errors.New("unreachable retry loop")
}
指数退避加抖动是为了避免大量客户端在同一时刻同时重试,把瞬时故障放大成雪崩。实际工程中还可以加入熔断机制,例如连续失败超过阈值后,短暂熔断所有对该推理服务的请求,等待一段时间再放行少量探测请求。熔断可以借助简单的滑动窗口计数器实现,不必引入复杂框架。
另一个容易被忽略的是请求体的复用。高并发下如果每次重试都重新构造请求体,会造成不必要的内存分配。如果使用bytes.Reader作为请求体,可以通过req.GetBody让http.Client在重试时重新读取。或者在重试循环中每次重新设置req.Body。对于JSON请求,更推荐使用[]byte配合bytes.NewReader,并在重试前重新创建请求,避免请求体被消耗后无法再次发送。
流式推理响应处理与背压控制
很多推理API支持SSE格式的流式输出,服务端会持续推送data:行,直到结束标记。如果SDK只是简单地用io.ReadAll读取完整响应再解析,就完全丧失了流式输出降低首字延迟的优势。高并发场景下,必须逐行读取并尽快把token回调给调用方。Go的bufio.Scanner可以按行读取,但默认缓冲区只有64KB,一条超长的data行会触发token too long错误。对于可能返回大段JSON的推理API,需要提前调大Scanner的缓冲区容量。
流式处理的难点在于背压控制。服务端推送速度可能远快于调用方消费速度,如果SDK内部无限制地读取并缓存,内存会迅速膨胀。常见的做法是使用带缓冲的channel作为中间层,并设置一个合理的缓冲大小。当channel写满时,读取goroutine阻塞,这样会反向给服务端施加背压,降低推送速度。下面代码实现了一个基于channel的流式token转发模型。
package inference
import (
"bufio"
"context"
"net/http"
"strings"
)
type StreamEvent struct {
Data string
Error error
}
func (c *Client) streamCompletion(ctx context.Context, req *http.Request) (<-chan StreamEvent, error) {
resp, err := c.httpClient.Do(req)
if err != nil {
return nil, err
}
events := make(chan StreamEvent, 64)
go func() {
defer close(events)
defer resp.Body.Close()
scanner := bufio.NewScanner(resp.Body)
scanner.Buffer(make([]byte, 64*1024), 1024*1024)
for scanner.Scan() {
line := scanner.Text()
if !strings.HasPrefix(line, "data:") {
continue
}
payload := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
if payload == "[DONE]" {
return
}
select {
case events <- StreamEvent{Data: payload}:
case <-ctx.Done():
return
}
}
if err := scanner.Err(); err != nil {
select {
case events <- StreamEvent{Error: err}:
case <-ctx.Done():
}
}
}()
return events, nil
}
上面这段代码中,events channel的缓冲大小设置为64,意味着最多允许64个待消费的数据块。如果调用方处理速度跟不上,读取goroutine会阻塞在发送操作上,这样resp.Body的读取自然停止,服务端也会感知到TCP窗口减小,从而降低推送速率。这种背压机制比无限制读取到内存要安全得多。
调用方在遍历channel时,需要区分正常数据、结束和错误三种情况。还可以传入一个独立的context,当调用方主动取消时,读取goroutine立即退出并关闭连接。注意这里关闭channel必须由读取goroutine完成,避免向已关闭的channel发送数据导致panic。如果希望支持多路复用,可以在SDK内部实现一个事件路由器,但高并发下单连接单goroutine的模式通常更高效。
可观测性指标与实际压测调优
没有可观测性的高并发SDK就像在黑夜里开车。至少需要记录三类指标:请求延迟分布、错误率、连接复用情况。Go标准库提供了expvar,可以方便地把内存中的计数器暴露出来。更完善的方案是接入Prometheus的client_golang,用Histogram记录延迟,用Counter记录错误和总请求数。连接复用率可以通过Transport的接口间接统计,例如在每次请求前后检查resp.Proto和连接数量,或者直接使用httptrace的GotConnInfo来记录是否复用了连接。
实际压测时,常常会发现两个奇怪的现象。第一个是连接数明明够用,但吞吐上不去,这多半是因为每个请求的goroutine在等待响应时占用大量栈内存,导致GC压力上升。第二个是延迟偶尔出现毛刺,排查后发现是TLS握手没有复用会话,或者DNS解析被频繁触发。针对前一个问题,可以适当降低单机并发上限,用队列削峰填谷;针对后一个问题,可以缓存DNS结果或者使用长连接。
压测工具的并发模型也会影响结果。如果使用简单的for循环加goroutine无上限发起请求,测试早于服务端崩溃。更合理的方式是使用固定并发数的worker池,每个worker串行发送请求,观察P50、P95、P99延迟和错误率。经过调优后,一个封装良好的Go SDK在单机8核环境下,对上万QPS的短推理请求通常可以保持P95延迟在200毫秒以内,连接复用率稳定在90%以上。这些数字需要根据具体推理服务调整,但设计原则是通用的:传输层要复用连接,请求层要有限重试,流式层要背压,观测层要暴露关键指标。