37DATA

Server-Sent Events (SSE)简单实践的一些问题和思考

序言

SSE是什么

WIKI的介绍如下(百度没有词条,看看维基的)

Server-Sent Events (SSE) is a server push technology enabling a client to receive automatic updates from a server via an HTTP connection, and describes how servers can initiate data transmission towards clients once an initial client connection has been established. They are commonly used to send message updates or continuous data streams to a browser client and designed to enhance native, cross-browser streaming through a JavaScript API called EventSource, through which a client requests a particular URL in order to receive an event stream. The EventSource API is standardized as part of HTML5 by the WHATWG. The media type for SSE is text/event-stream.

(From Wikipedia, the free encyclopedia)

网上很多人的blog都提到,SSE是HTML5的一个标准套件、是基于HTTP的一个websocket的轻量化替代协议、可以服务端流式传输到客户端、media类型是text/event-steam。
这些简介已经概括完了SSE的功能,但是少了WIKI里提到的重要的点:客户端通过通过称为 EventSource 的 JavaScript API增强本机跨浏览器流交互、请求特定URL接收事件流。所以,SSE本质上是基于HTTP + EventSource API做的一种服务端推送机制,被纳入了HTML5的标准里。

为什么会用到这个

为什么不是TCP/UDP或者websocket呢?大概因为SSE比起来更加轻量化吧。相比于前者需要自行管理连接的开发模式,基于HTTP的SSE能够让你以与传统短连接开发模式几乎一样的思路来开发,成本很低。
当然应用场景也有限制。SSE在收到请求后,便只能由服务器单向不断推送数据到客户端,而无法做到双工通信,因此应用场景也不如TCP/UDP或websocket。
但对于最近爆火的chatgpt来说,SSE的“单次请求、多段返回”正是一个极好的应用。

SSE的使用

一个简单实现:
package mainimport (    "context"    "fmt"    "log"    "net/http"    "time")// MessageEvent 事件type MessageEvent struct {    Id    string // 事件ID    Event string // 事件类型    Data  string // 发送数据}// 事件转字符串func (e MessageEvent) String() string {    return fmt.Sprintf("id:%s\nevent:%s\ndata:%s\n\n", e.Id, e.Event, e.Data)}func main() {    // 事件消息通道    messageChan := make(chan MessageEvent)    // HTTP 请求处理函数    http.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) {       // 设置 response header       // 媒体类型设为事件流       w.Header().Set("Content-Type", "text/event-stream")       // 因为是不断往缓冲区写入并flush出去,因此要禁用缓存       w.Header().Set("Cache-Control", "no-cache")       // 保持tcp连接,SSE是基于HTTP的,因此可以减少连接开销       w.Header().Set("Connection", "keep-alive")       // 设置响应码       w.WriteHeader(http.StatusOK)       // 结束上下文       t, cancel := context.WithCancel(context.Background())       // 写入事件信息       go func() {          for i := 0; i < 3; i++ {             messageChan <- MessageEvent{                Id:    fmt.Sprintf("%d", i),                Event: "test",                Data:  fmt.Sprintf("Message %d from server", i),             }             time.Sleep(time.Second)          }          // 结束上下文          cancel()       }()    L:       // 监听事件消息       for {          select {          case <-t.Done():             // 结束连接             log.Println("send finish")             break L          case message := <-messageChan:             // 发送数据             fmt.Fprintf(w, "%s", &message)             // 从缓冲区刷新数据出去             w.(http.Flusher).Flush()          }       }    })    // 启动监听    err := http.ListenAndServe(":8080", nil)    if err != nil {       panic(err)    }}
实现效果:Image

在技术中心tcf框架的实践

首先我们按上面的简单思路实现一个
package controllerimport (    "context"    "fmt"    v1 "ig-gptword-service/api/v1"    "log"    "net/http"    "time"    "github.com/gogf/gf/v2/frame/g")var (    TestSse = cTestSse{})// MessageEvent 事件type MessageEvent struct {    Id    string // 事件ID    Event string // 事件类型    Data  string // 发送数据}// 事件转字符串func (e MessageEvent) String() string {    return fmt.Sprintf("id:%s\nevent:%s\ndata:%s\n\n", e.Id, e.Event, e.Data)}type cTestSse struct{}func (c *cTestSse) Translate(ctx context.Context, req *v1.SseClientReq) (res *v1.SseClientRes, err error) {    r := g.RequestFromCtx(ctx)    w := r.Response.ResponseWriter.RawWriter()    // 事件消息通道    messageChan := make(chan MessageEvent)    // 设置 response header    // 媒体类型设为事件流    w.Header().Set("Content-Type", "text/event-stream")    // 因为是不断往缓冲区写入并flush出去,因此要禁用缓存    w.Header().Set("Cache-Control", "no-cache")    // 保持tcp连接,SSE是基于HTTP的,因此可以减少连接开销    w.Header().Set("Connection", "keep-alive")    // 设置响应码    w.WriteHeader(http.StatusOK)    // 结束上下文    t, cancel := context.WithCancel(context.Background())    // 向事件通道中写入数据    go func() {       for i := 0; i < 3; i++ {          messageChan <- MessageEvent{             Id:    fmt.Sprintf("%d", i),             Event: "test",             Data:  fmt.Sprintf("Message %d from server", i),          }          time.Sleep(time.Second)       }       // 结束上下文       cancel()    }()L:    // 循环监听 SSE 事件通道    for {       select {       case <-t.Done():          // 结束连接          log.Println("send finish")          break L       case message := <-messageChan:          // 向客户端发送 SSE 事件          fmt.Fprintf(w, "%s", &message)          // 刷新 response buffer          w.(http.Flusher).Flush()       }    }    return}
在浏览器中可以获得相同的成功结果,但是在服务端却出现了报错

Image

Image

原因是,GF没有对SSE的封装,都是在整个controller结束后统一返回(最后flush的时候会设置状态码),而SSE需要提前返回并设置状态码(即 w.WriteHeader(http.StatusOK) ),而header头的状态码仅允许被设置一次,因此出现报错。
因此,可以考虑在return前对连接进行hijack,使其脱离框架的请求生命周期管理,交由我们自行管理
// 劫持原连接,接管后续生命周期conn, _, err := r.Response.ResponseWriter.Hijack()if err != nil {    log.Println("close sse connection failed: ", err.Error())    return}// 关闭连接err = conn.Close()if err != nil {    log.Println("close sse connection failed: ", err.Error())}return
这时候,我们发现该错误消失了,但同时出现了另一个细节问题:数据明明发送成功了,但请求的状态却显示失败。

Image

Image

进去看看header头,发现在未指定Transfer-Encoding头时,会被自动指定chunked。

Image

在SSE里,该次HTTP请求的content-length无法提前预知,因为每一段都是独立返回的,因此无法通过返回的内容长度来判定请求是否结束,因此chunked用来提示客户端,需要通过其他方式来判断。
过去由于有框架存在,这样的细节我们无需关注,可是我们截获请求后,这部分就需要我们自行处理了。因此我们需要在最后断开连接前给出结束标志:0\r\n\r\n
// 劫持原连接,接管后续生命周期conn, _, err := r.Response.ResponseWriter.Hijack()if err != nil {    log.Println("close sse connection failed: ", err.Error())    return}// 关闭连接conn.Write([]byte("0\r\n\r\n"))err = conn.Close()if err != nil {    log.Println("close sse connection failed: ", err.Error())}return
成功
Image

一点应用场景的思考

  • 报表应用
    每个报表查询背后是无数个SQL,这些SQL执行效率各有差异,以往短连接开发模式有水桶效应,报表呈现效率取决于最慢的SQL。
    如果能在报表应用上使用SSE,可以打破水桶效应,实时反馈每个SQL、指标的查询进度和结果。
    而使用SSE相比于websocket或TCP/UDP,对现有应用、开发人员来说压力最小,改造成本最低。
  • 订阅服务
  • 取代轮询