跳到主要内容

实时通信之 SSE

前言

Server send event,简称SSE,是一种实时通信技术,作用类似于WebSocket,但SSE是单向的,WebSocket是双向的。某些场景下,单向通信更加高效。

在只需要服务器发送消息给客户端的场景中,比如消息推送,使用SSE技术会更加合适。另外,SSE技术是使用HTTP协议传输的,因此无需额外的实现就可以使用。此外,SSE还有一些特殊功能,比如自动重连接、Event IDs、发送自定义事件等。

SSE可以看作是一个客户端去从服务器端订阅一条事件流,之后服务端可以发送订阅特定事件的消息给客户端,直到服务端或者客户端关闭该连接。

本质上,SSE通过一个独立的Ajax请求从客户端向服务端传送数据。如果对交互的实时性要求很高,可以考虑使用websocket。

示例代码

Go

  • sse-go/sse/broker.go
package sse

import (
"fmt"
"log"
"net/http"
)

type Broker struct {
Notifier chan []byte // Events 推送到此信道
newClients chan chan []byte // 客户端连接
closingClients chan chan []byte // 关闭客户端连接
clients map[chan []byte]bool // 客户端连接注册
}

func (broker *Broker) listen() {
for {
select {
case s := <-broker.newClients:
// 新连接时注册
broker.clients[s] = true
log.Printf("New client added, total: %d\n", len(broker.clients))
case s := <-broker.closingClients:
// 新连接关闭时移除注册
delete(broker.clients, s)
log.Printf("Removed client, total: %d\n", len(broker.clients))
case event := <-broker.Notifier:
// 广播消息
for clientMessageChan := range broker.clients {
clientMessageChan <- event
}
}
}
}

// 实现 http.Handler 接口
func (broker *Broker) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
flusher, ok := rw.(http.Flusher)
if !ok {
http.Error(rw, "Streaming unsupported!", http.StatusInternalServerError)
return
}

// 设置响应头SSE
rw.Header().Set("Content-Type", "text/event-stream")
rw.Header().Set("Cache-Control", "no-cache")
rw.Header().Set("Connection", "keep-alive")
rw.Header().Set("Access-Control-Allow-Origin", "*")

messageChan := make(chan []byte) // 每一个连接都有自己的消息信道
// 新连接通知broker
broker.newClients <- messageChan
defer func() {
broker.closingClients <- messageChan
}()

ctx := req.Context()
go func() {
<-ctx.Done()
broker.closingClients <- messageChan
}()

// 阻塞等待
for {
fmt.Fprintf(rw, "data: %s\n\n", <-messageChan)
flusher.Flush()
}
}

  • sse-go/sse/sse.go
package sse


func NewSSE() (broker *Broker) {
broker = &Broker{
Notifier: make(chan []byte, 1),
newClients: make(chan chan []byte),
closingClients: make(chan chan []byte),
clients: make(map[chan []byte]bool),
}
// 开启监听
go broker.listen()
return
}
  • sse-go/main.go
package main

import (
"fmt"
"log"
"net/http"
"sse-go/sse"
"text/template"
"time"
)

func main() {
broker := sse.NewSSE()
go func() {
for {
time.Sleep(time.Second * 2)
data := fmt.Sprintf("==>%s", time.Now().Format("2006-01-02 15:04:05"))
log.Println("Sending event data")
broker.Notifier <- []byte(data)
}
}()

go func() {
for {
time.Sleep(time.Second * 1)
if time.Now().Second() % 2 == 0 {
data := fmt.Sprintf("-->%s", time.Now().Format("2006-01-02 15:04:05"))
log.Println("Sending event data2")
broker.Notifier <- []byte(data)
}
}
}()

http.HandleFunc("/", handleIndex)
http.Handle("/sse", broker)
log.Fatal("error: ", http.ListenAndServe(":8000", nil))
}

func handleIndex(writer http.ResponseWriter, request *http.Request) {
t, _ := template.ParseFiles("index.html")
t.Execute(writer, nil)
}
  • sse-go/index.html
<!DOCTYPE html>
<html lang="en">
<head>
<title>Go SSE</title>
</head>
<h1>Go SSE</h1>
<div id="msg"></div>
<body>
<script>
window.onload = function() {
var client = new EventSource("http://localhost:8000/sse")
client.onmessage = function(msg) {
document.getElementById("msg").innerHTML += msg.data + "<br/>"
console.log(msg)
}
}
</script>
</body>
</html>

Python-Flask

  • server.py
# 导入所需的模块
import json
import time
import datetime
from flask import Flask, request, Response, render_template
# from queue import Queue

app = Flask(__name__)

# event_queue = Queue(maxsize=100)

# 解决跨域问题
@app.after_request
def after_request(response):
response.headers.add('Access-Control-Allow-Origin', '*')
response.headers.add('Access-Control-Allow-Methods', 'GET,PUT,POST,DELETE,OPTIONS')
response.headers.add('Access-Control-Allow-Headers', 'Content-Type,Authorization')
response.headers.add('Access-Control-Allow-Credentials', 'true')
return response

# 获取当前时间,并转换为 JSON 格式
def get_time_json():
dt_ms = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')
return json.dumps({'time': dt_ms}, ensure_ascii=False)

# 设置路由,返回 SSE 流
@app.route('/')
def hello_world():
return render_template('sse.html')

@app.route('/sse')
def stream():
user_id = request.args.get('user_id') # 可选,用于区分不同用户的连接
print(user_id)

def eventStream():
id = 0
while True:
id += 1
time.sleep(.1) # 每秒发送约 1 条数据
event_name = 'time_reading'
str_out = f'id: {id}\nevent: {event_name}\ndata: {get_time_json()}\n\n'
# print(str_out) # 在服务器端打印发送的数据
yield str_out

return Response(eventStream(), mimetype="text/event-stream")

if __name__ == '__main__':
app.run(host='0.0.0.0', port=8000, debug=True)
  • templates/sse.html
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Server-Sent Events</title>
</head>
<body>
<div id="content">
<h1>Server-Sent Events</h1>
<p>time: <span id="time_show"></span></p>
</div>
</body>
<script>
function connectSSE() {
if (window.EventSource) {
let sse_url = 'http://localhost:8000/sse';
// 创建 EventSource 对象连接服务器
const source = new EventSource(sse_url);

// 连接成功后会触发 open 事件
source.addEventListener('open', () => {
console.log('Connected');
}, false);

// 当接收到服务器端发送的数据时会触发 time_reading 事件
source.addEventListener('time_reading', function (e) {
console.log("time_reading", e.data);
// 把时间显示到页面上,这里需要解析一下 JSON
document.getElementById("time_show").innerHTML = JSON.parse(e.data).time;
}, false);

// 服务器发送信息到客户端时,如果没有 event 字段,默认会触发 message 事件
source.addEventListener('message', e => {
console.log(`data: ${e.data}`);
}, false);

// 连接异常时会触发 error 事件并自动重连
source.addEventListener('error', e => {
if (e.target.readyState === EventSource.CLOSED) {
console.log('Disconnected');
} else if (e.target.readyState === EventSource.CONNECTING) {
console.log('Connecting...');
}
}, false);
} else {
console.error('Your browser doesn\'t support SSE');
}
}

connectSSE();
</script>
</html>

参考

  • 汪明《Go并发编程实战》