SSE åè®®æ¯æ
æ¬æä»ç»å¦ä½å¨ Web äºå½æ°ä¸å®ç° SSEï¼Server-Sent Eventsï¼æå¡ç«¯æ¨éï¼æ¯ææå¡ç«¯å客æ·ç«¯å忍é宿¶æ°æ®æµã
ä»ä¹æ¯ SSEâ
ãSSEï¼Server-Sent Eventsï¼ãæ¯ä¸ç§æå¡å¨å客æ·ç«¯æ¨é宿¶æ°æ®çææ¯ï¼åºäº HTTP åè®®å®ç°ãä¸ WebSocket ä¸åï¼SSE æ¯ååéä¿¡ï¼åªæ¯ææå¡ç«¯å客æ·ç«¯æ¨éæ°æ®ã
主è¦ç¹ç¹ï¼
- å忍éï¼æå¡ç«¯ä¸»å¨å客æ·ç«¯æ¨éæ°æ®ï¼å®¢æ·ç«¯åªè½æ¥æ¶
- åºäº HTTPï¼ä½¿ ç¨æ å HTTP åè®®ï¼å ¼å®¹æ§å¥½ï¼æäºå®ç°
- èªå¨éè¿ï¼å®¢æ·ç«¯æçº¿åä¼èªå¨éæ°è¿æ¥
- ææ¬æ ¼å¼ï¼ä¼ è¾çæ°æ®ä¸ºææ¬æ ¼å¼ï¼éå¸¸æ¯ JSONï¼
- è½»éå®ç°ï¼ç¸æ¯ WebSocketï¼å®ç°æ´ç®åï¼èµæºå ç¨æ´å°
å ¸ååºç¨åºæ¯ï¼
- AI å¯¹è¯æµå¼è¾åº
- 宿¶æ¥å¿æ¨é
- è¿åº¦æ´æ°éç¥
- æå¡ç«¯ç¶æçæ§
- 宿¶æ°æ®çæ¿
å·¥ä½åçâ
åè®®å¯ç¨â
SSE åè®®å¨ Web äºå½æ°ä¸é»è®¤æ¯æï¼æ é卿§å¶å°è¿è¡ä»»ä½é¢å¤é ç½®å³å¯ä½¿ç¨ã
建ç«è¿æ¥â
- 客æ·ç«¯éè¿æ å HTTP 请æ±å»ºç« SSE è¿æ¥
- æå¡ç«¯è¿å
Content-Type: text/event-streamååºå¤´ - è¿æ¥ä¿ææå¼ç¶æï¼æå¡ç«¯æç»æ¨éæ°æ®
è¿æ¥çå½å¨æâ
- è°ç¨å¯¹åºå
³ç³»ï¼ä¸æ¬¡ SSE è¿æ¥ççå½å¨æçåäºä¸æ¬¡å½æ°è°ç¨è¯·æ±
- è¿æ¥å»ºç« = 请æ±åèµ·
- è¿æ¥æå¼ = 请æ±ç»æ
- å®ä¾æ å°ï¼å½æ°å®ä¾ä¸ SSE è¿æ¥æ¯ä¸ä¸å¯¹åºçï¼åä¸å®ä¾å¨æä¸æ¶å»ä» å¤çä¸ä¸ª SSE è¿æ¥ï¼æ°è¿æ¥ä¼å¯å¨æ°çå®ä¾
- è¿æ¥ä¿æï¼è¿æ¥å»ºç«åï¼å®ä¾æç»è¿è¡ï¼éè¿æµå¼ååºæ¨éæ°æ®
- è¿æ¥ç»æï¼å½ SSE è¿æ¥æå¼ææå¡ç«¯è°ç¨
end()æ¶ï¼å¯¹åºç彿°å®ä¾åæ¢è¿è¡
使ç¨éå¶â
å¨ä½¿ç¨ SSE æ¶ï¼éè¦æ³¨æä»¥ä¸éå¶ï¼
| éå¶é¡¹ | 说æ |
|---|---|
| æ§è¡è¶ æ¶æ¶é´ | è¿æ¥æç»æ¶é´å彿°æå¤§è¿è¡æ¶é¿éå¶ |
| å¹¶åè¿æ¥ | æ¯ä¸ªè¿æ¥å¯å¨ç¬ç«å®ä¾ï¼åè´¦æ·å¹¶åé é¢éå¶ |
| æµè§å¨è¿æ¥æ°éå¶ | åä¸åå䏿µè§å¨å¯¹ SSE è¿æ¥æ°æéå¶ï¼é常为 6 ä¸ªï¼ |
| æ°æ®æ ¼å¼ | åªè½ä¼ è¾ææ¬æ°æ®ï¼å¤ææ°æ®éåºåå为 JSON |
æä½æ¥éª¤â
æ¥éª¤1ï¼ç¼åæå¡ç«¯ä»£ç â
SSE åè®®é»è®¤æ¯æï¼æ é卿§å¶å°å¼å¯ãæ ¹æ®æ¨ä½¿ç¨çç¼ç¨è¯è¨åæ¡æ¶ ï¼ç¼å SSE æå¡ç«¯ä»£ç ã
- Node.js (Express)
- Node.js (Koa)
- Python (Flask)
- Python (FastAPI)
å®è£ ä¾èµ
å®è£ Express æ¡æ¶ï¼
npm install express
å¨ package.json 䏿·»å ä¾èµï¼
{
"dependencies": {
"express": "^4.18.0"
}
}
ç¼åæå¡ç«¯ä»£ç
ä½¿ç¨ Express å®ç° SSE æµå¼ååºï¼
const express = require('express');
const app = express();
// SSE è·¯ç±
app.get('/stream', (req, res) => {
// 设置 SSE ååºå¤´
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
console.log('æ°ç SSE è¿æ¥å»ºç«');
// æµå¼æ¨éæ°æ®
const msg = ['SSE', 'empowering', 'GPT', 'applications', '!', 'Happy', 'chatting', '!'];
let index = 0;
const intervalId = setInterval(() => {
if (index < msg.length) {
// åé SSE æ¶æ¯
const data = {
id: index,
content: msg[index],
};
res.write(`data: ${JSON.stringify(data)}\n\n`);
index++;
} else {
// 宿æ¨éï¼å
³éè¿æ¥
clearInterval(intervalId);
res.end();
}
}, 1000);
// 客æ·ç«¯æå¼è¿æ¥æ¶æ¸
ç
req.on('close', () => {
clearInterval(intervalId);
console.log('SSE è¿æ¥å·²å
³é');
});
});
// å¯å¨æå¡ï¼çå¬ 9000 端å£
app.listen(9000, () => {
console.log('SSE æå¡å·²å¯å¨ï¼çå¬ç«¯å£ 9000');
});
å®è£ ä¾èµ
å®è£ Koa æ¡æ¶ï¼
npm install koa koa-router
ç¼åæå¡ç«¯ä»£ç
ä½¿ç¨ Koa å®ç° SSE æµå¼ååºï¼
const Koa = require('koa');
const Router = require('koa-router');
const app = new Koa();
const router = new Router();
router.get('/stream', async (ctx) => {
// 设置 SSE ååºå¤´
ctx.set({
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
});
console.log('æ°ç SSE è¿æ¥å»ºç«');
const msg = ['SSE', 'empowering', 'GPT', 'applications', '!', 'Happy', 'chatting', '!'];
// å建å¯è¯»æµ
const stream = new require('stream').PassThrough();
ctx.body = stream;
let index = 0;
const intervalId = setInterval(() => {
if (index < msg.length) {
const data = {
id: index,
content: msg[index],
};
stream.write(`data: ${JSON.stringify(data)}\n\n`);
index++;
} else {
clearInterval(intervalId);
stream.end();
}
}, 1000);
// 客æ·ç«¯æå¼è¿æ¥æ¶æ¸
ç
ctx.req.on('close', () => {
clearInterval(intervalId);
console.log('SSE è¿æ¥å·²å
³é');
});
});
app.use(router.routes()).use(router.allowedMethods());
app.listen(9000, () => {
console.log('SSE æå¡å·²å¯å¨ï¼çå¬ç«¯å£ 9000');
});
å®è£ ä¾èµ
å¨ requirements.txt 䏿·»å Flaskï¼
Flask
ç¼åæå¡ç«¯ä»£ç
ä½¿ç¨ Flask å®ç° SSE æµå¼ååºï¼
import json
import time
from flask import Flask, Response, stream_with_context
app = Flask(__name__)
@app.route('/stream')
def stream_data():
"""SSE æµå¼æ¨é"""
msg = ['SSE', 'empowering', 'GPT', 'applications', '!', 'Happy', 'chatting', '!']
def generate_response_data():
for i, word in enumerate(msg):
json_data = json.dumps({'id': i, 'content': word})
yield f"data: {json_data}\n\n"
time.sleep(1)
return Response(
stream_with_context(generate_response_data()),
mimetype="text/event-stream",
headers={
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
}
)
if __name__ == '__main__':
app.run(host='0.0.0.0', port=9000)
å®è£ ä¾èµ
å¨ requirements.txt 䏿·»å FastAPI å uvicornï¼
fastapi
uvicorn
ç¼åæå¡ç«¯ä»£ç
ä½¿ç¨ FastAPI å®ç° SSE æµå¼ååºï¼
import json
import asyncio
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
@app.get('/stream')
async def stream_data():
"""SSE æµå¼æ¨é"""
async def generate_response_data():
msg = ['SSE', 'empowering', 'GPT', 'applications', '!', 'Happy', 'chatting', '!']
for i, word in enumerate(msg):
json_data = json.dumps({'id': i, 'content': word})
yield f"data: {json_data}\n\n"
await asyncio.sleep(1)
return StreamingResponse(
generate_response_data(),
media_type="text/event-stream",
headers={
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
}
)
if __name__ == '__main__':
import uvicorn
uvicorn.run(app, host='0.0.0.0', port=9000)