
FastAPI实时通信SSE实战
大约 7 分钟
FastAPI实时通信SSE实战
如果说WebSocket是实时通信的"双向高速公路",那么SSE(Server-Sent Events)就是"单向快车道"——专门用于服务器向客户端推送数据。对于测试平台来说,SSE就像是"实时监控大屏",能够将测试执行状态、系统监控数据实时推送给前端,让测试工程师能够实时掌握系统状态。
一、SSE简介:服务器推送的艺术
什么是SSE?
Server-Sent Events(SSE)是HTML5标准的一部分,它允许服务器主动向客户端推送数据:
SSE的特点:
- 单向通信:只能服务器向客户端推送
- 自动重连:连接断开后自动重连
- 简单易用:基于HTTP协议,无需特殊协议
- 事件驱动:支持自定义事件类型
- 轻量级:比WebSocket更轻量
SSE vs WebSocket vs 轮询
| 特性 | SSE | WebSocket | 轮询 |
|---|---|---|---|
| 通信方向 | 单向(服务器→客户端) | 双向 | 双向 |
| 协议 | HTTP | WebSocket | HTTP |
| 复杂度 | 简单 | 中等 | 简单 |
| 自动重连 | 是 | 否 | 否 |
| 实时性 | 高 | 最高 | 低 |
| 服务器压力 | 中 | 中 | 高 |
测试平台中的SSE应用场景
# SSE在测试平台中的典型应用
sse_use_cases = {
"测试执行监控": {
"场景": "实时推送测试用例执行状态",
"数据": "测试进度、结果、错误信息",
"频率": "按事件触发"
},
"系统性能监控": {
"场景": "实时推送系统性能指标",
"数据": "CPU、内存、网络使用率",
"频率": "每秒更新"
},
"日志流式输出": {
"场景": "实时显示应用日志",
"数据": "错误日志、访问日志、调试信息",
"频率": "实时流式"
},
"通知推送": {
"场景": "系统通知和告警",
"数据": "测试完成通知、系统告警",
"频率": "按事件触发"
}
}二、FastAPI中的SSE实现
基础SSE实现
# sse_basic.py
from fastapi import FastAPI, Response
from fastapi.responses import StreamingResponse
import asyncio
import json
import time
from datetime import datetime
from typing import AsyncGenerator
app = FastAPI(title="测试平台SSE服务")
# 基础SSE响应
@app.get("/sse/basic")
async def basic_sse():
"""基础SSE示例"""
async def event_generator() -> AsyncGenerator[str, None]:
counter = 0
while True:
# SSE格式:data: {content}\n\n
yield f"data: 当前时间: {datetime.now()}, 计数: {counter}\n\n"
counter += 1
await asyncio.sleep(1)
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*",
"Access-Control-Allow-Headers": "Cache-Control"
}
)
# 带事件类型的SSE
@app.get("/sse/events")
async def event_sse():
"""带事件类型的SSE"""
async def event_generator() -> AsyncGenerator[str, None]:
events = ["info", "warning", "error", "success"]
counter = 0
while True:
event_type = events[counter % len(events)]
data = {
"id": counter,
"message": f"这是一个{event_type}事件",
"timestamp": datetime.now().isoformat(),
"level": event_type
}
# 带事件类型的SSE格式
yield f"event: {event_type}\n"
yield f"data: {json.dumps(data)}\n\n"
counter += 1
await asyncio.sleep(2)
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*"
}
)
# 带ID的SSE(支持断线重连)
@app.get("/sse/with-id")
async def sse_with_id(last_event_id: str = None):
"""带ID的SSE,支持断线重连"""
async def event_generator() -> AsyncGenerator[str, None]:
# 从last_event_id开始发送
start_id = int(last_event_id) + 1 if last_event_id else 0
for event_id in range(start_id, start_id + 100):
data = {
"event_id": event_id,
"message": f"事件 {event_id}",
"timestamp": datetime.now().isoformat()
}
yield f"id: {event_id}\n"
yield f"data: {json.dumps(data)}\n\n"
await asyncio.sleep(1)
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*"
}
)测试执行状态推送
# test_execution_sse.py
from fastapi import FastAPI, BackgroundTasks, HTTPException
import asyncio
import json
from datetime import datetime
from typing import Dict, List, AsyncGenerator
from enum import Enum
app = FastAPI()
class TestStatus(str, Enum):
PENDING = "pending"
RUNNING = "running"
PASSED = "passed"
FAILED = "failed"
SKIPPED = "skipped"
# 全局状态存储
test_execution_status: Dict[int, dict] = {}
active_connections: List[AsyncGenerator] = []
class TestExecutionManager:
"""测试执行管理器"""
@staticmethod
async def execute_test_case(test_case_id: int):
"""模拟测试用例执行"""
# 更新状态为运行中
test_execution_status[test_case_id] = {
"id": test_case_id,
"status": TestStatus.RUNNING,
"progress": 0,
"message": "测试开始执行",
"start_time": datetime.now().isoformat()
}
# 模拟测试执行过程
for progress in range(0, 101, 10):
test_execution_status[test_case_id].update({
"progress": progress,
"message": f"测试执行中... {progress}%"
})
await asyncio.sleep(0.5) # 模拟执行时间
# 随机决定测试结果
import random
success = random.choice([True, True, False]) # 70%成功率
test_execution_status[test_case_id].update({
"status": TestStatus.PASSED if success else TestStatus.FAILED,
"progress": 100,
"message": "测试执行完成" if success else "测试执行失败",
"end_time": datetime.now().isoformat(),
"result": {
"success": success,
"response_time": random.randint(100, 2000),
"error_message": None if success else "断言失败:期望状态码200,实际404"
}
})
@app.post("/api/testcases/{test_case_id}/execute")
async def execute_test_case(test_case_id: int, background_tasks: BackgroundTasks):
"""执行测试用例"""
# 检查是否已在执行
if test_case_id in test_execution_status:
current_status = test_execution_status[test_case_id]["status"]
if current_status == TestStatus.RUNNING:
raise HTTPException(status_code=400, detail="测试用例正在执行中")
# 添加后台任务
background_tasks.add_task(TestExecutionManager.execute_test_case, test_case_id)
return {
"message": f"测试用例 {test_case_id} 已开始执行",
"test_case_id": test_case_id,
"status": "submitted"
}
@app.get("/sse/test-execution")
async def test_execution_stream():
"""测试执行状态SSE流"""
async def event_generator() -> AsyncGenerator[str, None]:
last_sent_status = {}
try:
while True:
# 检查所有测试用例状态变化
for test_case_id, current_status in test_execution_status.items():
# 只发送状态有变化的测试用例
if (test_case_id not in last_sent_status or
last_sent_status[test_case_id] != current_status):
yield f"event: test_status_update\n"
yield f"data: {json.dumps(current_status)}\n\n"
last_sent_status[test_case_id] = current_status.copy()
await asyncio.sleep(0.1) # 100ms检查一次
except asyncio.CancelledError:
print("客户端断开连接")
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*"
}
)
@app.get("/api/testcases/{test_case_id}/status")
async def get_test_status(test_case_id: int):
"""获取测试用例状态"""
if test_case_id not in test_execution_status:
raise HTTPException(status_code=404, detail="测试用例状态不存在")
return test_execution_status[test_case_id]系统监控数据推送
# system_monitor_sse.py
import psutil
import asyncio
import json
from datetime import datetime
from typing import AsyncGenerator
@app.get("/sse/system-monitor")
async def system_monitor_stream():
"""系统监控数据SSE流"""
async def monitor_generator() -> AsyncGenerator[str, None]:
while True:
try:
# 获取系统性能数据
cpu_percent = psutil.cpu_percent(interval=None)
memory = psutil.virtual_memory()
disk = psutil.disk_usage('/')
network = psutil.net_io_counters()
# 构造监控数据
monitor_data = {
"timestamp": datetime.now().isoformat(),
"cpu": {
"usage_percent": cpu_percent,
"count": psutil.cpu_count()
},
"memory": {
"total": memory.total,
"available": memory.available,
"used": memory.used,
"usage_percent": memory.percent
},
"disk": {
"total": disk.total,
"used": disk.used,
"free": disk.free,
"usage_percent": (disk.used / disk.total) * 100
},
"network": {
"bytes_sent": network.bytes_sent,
"bytes_recv": network.bytes_recv,
"packets_sent": network.packets_sent,
"packets_recv": network.packets_recv
}
}
yield f"event: system_monitor\n"
yield f"data: {json.dumps(monitor_data)}\n\n"
await asyncio.sleep(1) # 每秒更新一次
except Exception as e:
error_data = {
"timestamp": datetime.now().isoformat(),
"error": str(e),
"type": "monitor_error"
}
yield f"event: error\n"
yield f"data: {json.dumps(error_data)}\n\n"
await asyncio.sleep(5) # 错误时等待5秒再重试
return StreamingResponse(
monitor_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*"
}
)
# 日志流式输出
@app.get("/sse/logs")
async def log_stream(level: str = "INFO"):
"""日志流式输出"""
async def log_generator() -> AsyncGenerator[str, None]:
import logging
import random
log_levels = ["DEBUG", "INFO", "WARNING", "ERROR"]
log_messages = [
"用户登录成功",
"测试用例执行完成",
"数据库连接异常",
"API请求超时",
"缓存更新成功",
"文件上传失败"
]
while True:
# 生成模拟日志
log_level = random.choice(log_levels)
message = random.choice(log_messages)
# 只推送指定级别及以上的日志
if log_levels.index(log_level) >= log_levels.index(level):
log_data = {
"timestamp": datetime.now().isoformat(),
"level": log_level,
"message": message,
"module": "test_platform",
"thread": "main"
}
yield f"event: log\n"
yield f"data: {json.dumps(log_data)}\n\n"
await asyncio.sleep(random.uniform(0.5, 3)) # 随机间隔
return StreamingResponse(
log_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*"
}
)三、前端React集成
基础EventSource使用
// src/hooks/useSSE.js
import { useState, useEffect, useRef } from 'react';
export function useSSE(url, options = {}) {
const [data, setData] = useState(null);
const [connectionStatus, setConnectionStatus] = useState('disconnected');
const [error, setError] = useState(null);
const eventSourceRef = useRef(null);
const {
onMessage,
onError,
onOpen,
onClose,
autoReconnect = true,
reconnectInterval = 3000
} = options;
useEffect(() => {
if (!url) return;
const connect = () => {
try {
setConnectionStatus('connecting');
eventSourceRef.current = new EventSource(url);
eventSourceRef.current.onopen = (event) => {
setConnectionStatus('connected');
setError(null);
onOpen?.(event);
};
eventSourceRef.current.onmessage = (event) => {
try {
const parsedData = JSON.parse(event.data);
setData(parsedData);
onMessage?.(parsedData, event);
} catch (e) {
setData(event.data);
onMessage?.(event.data, event);
}
};
eventSourceRef.current.onerror = (event) => {
setConnectionStatus('error');
setError('连接错误');
onError?.(event);
if (autoReconnect) {
setTimeout(connect, reconnectInterval);
}
};
} catch (err) {
setError(err.message);
setConnectionStatus('error');
}
};
connect();
return () => {
if (eventSourceRef.current) {
eventSourceRef.current.close();
setConnectionStatus('disconnected');
onClose?.();
}
};
}, [url]);
const disconnect = () => {
if (eventSourceRef.current) {
eventSourceRef.current.close();
setConnectionStatus('disconnected');
}
};
const reconnect = () => {
disconnect();
// 重新触发useEffect
setTimeout(() => {
if (eventSourceRef.current?.readyState === EventSource.CLOSED) {
connect();
}
}, 100);
};
return {
data,
connectionStatus,
error,
disconnect,
reconnect
};
}测试执行监控组件
// src/components/TestExecutionMonitor.jsx
import React, { useState, useEffect } from 'react';
import { useSSE } from '../hooks/useSSE';
function TestExecutionMonitor() {
const [testStatuses, setTestStatuses] = useState({});
const [logs, setLogs] = useState([]);
// 监听测试执行状态
const { connectionStatus: testConnectionStatus } = useSSE(
'http://localhost:8000/sse/test-execution',
{
onMessage: (data, event) => {
if (event.type === 'test_status_update') {
setTestStatuses(prev => ({
...prev,
[data.id]: data
}));
// 添加到日志
setLogs(prev => [...prev, {
id: Date.now(),
timestamp: new Date().toISOString(),
message: `测试用例 ${data.id}: ${data.message}`,
level: data.status === 'failed' ? 'error' : 'info'
}].slice(-100)); // 只保留最近100条
}
}
}
);
// 执行测试用例
const executeTestCase = async (testCaseId) => {
try {
const response = await fetch(`http://localhost:8000/api/testcases/${testCaseId}/execute`, {
method: 'POST'
});
if (response.ok) {
const result = await response.json();
console.log('测试用例提交成功:', result);
}
} catch (error) {
console.error('执行测试用例失败:', error);
}
};
const getStatusColor = (status) => {
const colors = {
pending: '#ffa500',
running: '#1890ff',
passed: '#52c41a',
failed: '#ff4d4f',
skipped: '#d9d9d9'
};
return colors[status] || '#d9d9d9';
};
const getProgressColor = (status) => {
return status === 'failed' ? '#ff4d4f' : '#1890ff';
};
return (
<div className="test-execution-monitor">
<div className="monitor-header">
<h2>测试执行监控</h2>
<div className="connection-status">
连接状态:
<span className={`status ${testConnectionStatus}`}>
{testConnectionStatus === 'connected' ? '已连接' :
testConnectionStatus === 'connecting' ? '连接中' : '已断开'}
</span>
</div>
</div>
<div className="test-controls">
<button onClick={() => executeTestCase(1)}>执行测试用例 1</button>
<button onClick={() => executeTestCase(2)}>执行测试用例 2</button>
<button onClick={() => executeTestCase(3)}>执行测试用例 3</button>
</div>
<div className="test-status-grid">
{Object.values(testStatuses).map(test => (
<div key={test.id} className="test-status-card">
<div className="test-header">
<h3>测试用例 {test.id}</h3>
<span
className="status-badge"
style={{ backgroundColor: getStatusColor(test.status) }}
>
{test.status}
</span>
</div>
<div className="test-progress">
<div className="progress-bar">
<div
className="progress-fill"
style={{
width: `${test.progress}%`,
backgroundColor: getProgressColor(test.status)
}}
/>
</div>
<span className="progress-text">{test.progress}%</span>
</div>
<div className="test-message">{test.message}</div>
{test.result && (
<div className="test-result">
<div>响应时间: {test.result.response_time}ms</div>
{test.result.error_message && (
<div className="error-message">
错误: {test.result.error_message}
</div>
)}
</div>
)}
<div className="test-times">
{test.start_time && (
<div>开始: {new Date(test.start_time).toLocaleTimeString()}</div>
)}
{test.end_time && (
<div>结束: {new Date(test.end_time).toLocaleTimeString()}</div>
)}
</div>
</div>
))}
</div>
<div className="logs-section">
<h3>执行日志</h3>
<div className="logs-container">
{logs.map(log => (
<div key={log.id} className={`log-entry ${log.level}`}>
<span className="log-time">
{new Date(log.timestamp).toLocaleTimeString()}
</span>
<span className="log-message">{log.message}</span>
</div>
))}
</div>
</div>
</div>
);
}
export default TestExecutionMonitor;