NDJSON 流式响应
流式是 fun 的一等能力:方法返回 *fun.Stream,响应变为 application/x-ndjson,每行一个 JSON。聊天、进度推送、AI 逐 token 输出都是这个模式。v1.3.6 补齐了长连接可靠性:服务端 25s 心跳 + 客户端原生重连。
两种签名
go
// 纯流式:所有消息都走流
func (s *ChatSvc) Chat(dto ChatDto) (*fun.Stream, error)
// 首条消息 + 后续流:T 作为第一条消息下发,后续走流
func (s *ChatSvc) Ask(dto AskDto) (string, *fun.Stream, error)(T, *Stream, error) 适合「先回元信息、再逐步推内容」的场景:例如先下发会话 id / 首屏数据,再推增量。
服务端写法
go
func (s *ChatSvc) Chat(dto ChatDto) (*fun.Stream, error) {
st := &fun.Stream{}
go func() {
for _, chunk := range chunks {
if err := st.Send(chunk); err != nil {
return // 连接断开,停止推送
}
}
st.Close() // 必须关闭,否则客户端一直等
}()
return st, nil
}三条纪律:
Send后检查 error——连接断开后返回错误,继续 Send 只会堆积,循环里要 return;- 结束时必须
Close()——不关流客户端会一直等(合法的零消息流也要 Close); - 推送逻辑放 goroutine——方法返回
st后框架才注入推送通道,Send在注入前会自动阻塞等待,不会竞态。
Stream API
| 方法 | 说明 |
|---|---|
Send(message any) error | 推送一条消息(自动 JSON 序列化,键转小写);流关闭 / 连接断开返回 error |
Close() | 结束流,触发 OnClose 回调;重复调用安全 |
OnClose(cb func()) | 注册清理回调(释放资源、上报指标);流已关闭时立即执行 |
go
st.OnClose(func() {
redis.Del("job:" + dto.Id) // 连接断开 / 正常结束都会触发
})错误处理
- 建流前的业务错误:正常
return nil, fun.Error(...),走普通Result错误响应(不会建立流); - 建流后的传输失败:
Send返回 error,业务侧自行收尾(OnClose里清理); - 框架在业务方法出错时会注入「已取消的流」,后续
Send立即返回错误,防止 goroutine 泄漏。
线协议
http
HTTP/1.1 200 OK
Content-Type: application/x-ndjson
{"status":0,"data":"你好"}
{"status":0,"data":",这是"}
{"status":0,"data":"流式示例"}每行一个完整 JSON(键同样递归转首字母小写);客户端逐行解析、逐行回调。
长连接保活:25s 心跳(v1.3.6)
流空闲时服务端每 25 秒发送一个空行心跳,穿透中间盒(移动 CGNAT 约 60s 超时、CDN 约 120s)防止连接被静默掐断:
- 心跳是空行,客户端解析器自动跳过,业务无感知;
- 配合服务端
SetTimeouts时注意idle要大于 25s(见服务器配置)。
客户端:原生重连(v1.3.6)
TS 客户端的流式调用支持指数退避重连与空闲看门狗:
ts
const c = api.create('/api')
// dto 传工厂函数:每次重连重新求值,用于刷新续传游标
c.chatSvc.chat(
() => ({ prompt: '讲个故事', cursor: getCursor() }),
(msg) => render(msg),
{
retry: { maxAttempts: 5, baseDelayMs: 500, maxDelayMs: 10000 },
idleTimeoutMs: 60_000, // 空闲看门狗:超时未收到任何字节(含心跳)判定流死亡并重连
}
)retry默认关闭(与 request 语义一致),显式传{...}开启;实际延迟为指数退避 × 0.75~1.25 抖动,maxAttempts默认 Infinity;idleTimeoutMs默认关闭;0 显式关闭;- 重连仅针对传输失败(status 4/5),业务错误不重试。
完整参数见 TS 客户端 · 流式调用。