Skip to content

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
}

三条纪律:

  1. Send 后检查 error——连接断开后返回错误,继续 Send 只会堆积,循环里要 return;
  2. 结束时必须 Close()——不关流客户端会一直等(合法的零消息流也要 Close);
  3. 推送逻辑放 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 客户端 · 流式调用

下一步

基于 MIT 许可发布