SSE 实时推送
EnerOS 提供 Server-Sent Events(SSE)端点,用于浏览器原生支持的单向实时事件推送。相比 WebSocket,SSE 更轻量、自动断线重连、走标准 HTTP 端口,适合告警看板、事件大屏、运维通知等只读订阅场景。SSE 基于 HTTP 长连接,由服务端通过 text/event-stream 流式推送事件,浏览器通过 EventSource API 原生接收。
SSE 协议说明
协议特性
SSE 是 HTML5 规范的一部分,基于 HTTP/1.1 长连接,由服务端持续推送文本事件流。主要特性:
- 单向通信:仅服务端→客户端,客户端无法通过同一连接发送数据
- 文本协议:仅支持 UTF-8 文本,不适合二进制数据
- 自动重连:浏览器原生实现断线重连,无需客户端代码
- 断点续传:通过
Last-Event-ID头补发遗漏事件 - HTTP 兼容:走标准 80/443 端口,穿透代理与防火墙
- 连接复用:HTTP/2 下可与其他请求复用同一 TCP 连接
协议格式
SSE 响应 Content-Type 为 text/event-stream,每条事件由若干字段行组成,以空行分隔:
event: <事件类型>
id: <事件 ID>
retry: <重连间隔毫秒>
data: <数据行 1>
data: <数据行 2>
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
event | string | 否 | 事件类型,客户端可按类型分发;缺省为 message |
id | string | 否 | 事件 ID,浏览器存储后重连时通过 Last-Event-ID 头发送 |
retry | int | 否 | 重连间隔(毫秒),浏览器将覆盖默认值 |
data | string | 是 | 事件数据,多行 data: 将以换行符拼接 |
以 : 开头的行为注释行,常用于心跳保活。
端点信息
| 项目 | 取值 |
|---|---|
| URL | GET http://localhost:8080/api/v1/events/stream |
| Content-Type | text/event-stream; charset=utf-8 |
| Cache-Control | no-cache |
| Connection | keep-alive |
| 认证 | Bearer Token(Header)或 Query 参数 ?token= |
| 心跳周期 | 15 秒(发送 :keepalive\n\n 注释行) |
| 重连间隔 | 默认 3 秒,可通过 retry 字段调整 |
| 单连接最大订阅 | 10 个主题 |
| 单租户最大连接 | 50 个并发 |
| 事件 payload 上限 | 64 KB |
| 事件保留时长 | 24 小时(用于断点续传) |
订阅参数
通过 Query 参数过滤事件流:
| 参数 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
topics | string | 否 | 全部 | 逗号分隔的主题列表,如 alarm,agent |
network_id | string | 否 | — | 限定电网 ID |
severity | string | 否 | — | 告警级别过滤:info/warn/error/critical |
agent_id | string | 否 | — | 限定 Agent ID |
since | datetime | 否 | — | 起始时间(RFC 3339) |
until | datetime | 否 | — | 结束时间(达到后自动关闭连接) |
last_event_id | string | 否 | — | 上次接收的事件 ID,用于断点续传 |
format | string | 否 | json | json / plain |
include_audit | bool | 否 | false | 是否包含审计事件(需审计权限) |
事件类型
SSE 端点推送以下事件类型,每个 event 字段对应一个分类:
| event | 说明 | payload 关键字段 | 触发频率 |
|---|---|---|---|
alarm | 告警产生/确认/清除 | alarm_id, severity, message | 实时 |
agent | Agent 状态变更 | agent_id, old_status, new_status | 实时 |
topology | 电网拓扑变更 | network_id, changes | 实时 |
powerflow | 潮流计算完成 | task_id, network_id, status | 任务驱动 |
constraint | 约束越限/恢复 | element, quantity, value, limit | 实时 |
device | 设备上下线 | device_id, status | 实时 |
task | 任务状态变更 | task_id, status | 任务驱动 |
audit | 审计事件(需权限) | actor, action, resource | 操作驱动 |
system | 系统通知(维护/升级) | type, message | 偶发 |
heartbeat | 心跳事件(应用层) | timestamp | 60 秒 |
消息格式示例
alarm 事件
event: alarm
id: evt_20260706_001
retry: 3000
data: {"id":"alm_001","severity":"critical","status":"active","type":"BRANCH_OVERLOAD","message":"Branch 1-2 overload","network_id":"net_8f3a2b","raised_at":"2026-07-06T08:30:00Z"}
agent 事件
event: agent
id: evt_20260706_002
data: {"agent_id":"dispatch_01","name":"DispatchAgent-01","old_status":"idle","new_status":"running","reason":"manual_start","timestamp":"2026-07-06T08:30:05Z"}
多行 data 示例
event: topology
id: evt_20260706_003
data: {"network_id":"net_8f3a2b","kind":"TopologyChanged"}
data: {"changes":[{"branch":"1-2","status":"open"},{"branch":"2-3","status":"closed"}]}
data: {"actor":{"type":"agent","id":"agent_self_healing_01"}}
心跳注释行
:keepalive
系统通知
event: system
id: evt_20260706_004
data: {"type":"maintenance_scheduled","message":"系统将于 30 分钟后进入维护模式","scheduled_at":"2026-07-06T09:00:00Z"}
浏览器示例
基础订阅
const es = new EventSource(
'http://localhost:8080/api/v1/events/stream?topics=alarm,agent&severity=critical',
{ withCredentials: true }
);
es.addEventListener('alarm', (e) => {
const data = JSON.parse(e.data);
console.log('告警:', data.severity, data.message);
showAlarmBanner(data);
});
es.addEventListener('agent', (e) => {
const data = JSON.parse(e.data);
console.log('Agent 状态:', data.agent_id, data.new_status);
updateAgentStatus(data.agent_id, data.new_status);
});
es.addEventListener('open', () => {
console.log('SSE 连接已建立');
});
es.addEventListener('error', (e) => {
if (es.readyState === EventSource.CLOSED) {
console.log('连接已关闭');
} else {
console.log('连接异常,浏览器将自动重连');
}
});
带 Token 订阅
浏览器 EventSource 不支持自定义 Header,需通过 Query 参数传递 Token:
const token = localStorage.getItem('eneros_token');
const es = new EventSource(
`http://localhost:8080/api/v1/events/stream?topics=alarm&token=${encodeURIComponent(token)}`,
{ withCredentials: true }
);
es.addEventListener('alarm', (e) => {
const data = JSON.parse(e.data);
console.log('告警:', data);
});
断点续传
服务端通过 Last-Event-ID 头自动处理断点续传。浏览器在重连时会自动附带此 Header,无需手动处理。但对于非浏览器客户端,需手动维护:
let lastEventId = null;
function connect() {
const url = lastEventId
? `http://localhost:8080/api/v1/events/stream?topics=alarm&last_event_id=${lastEventId}`
: 'http://localhost:8080/api/v1/events/stream?topics=alarm';
const es = new EventSource(url);
es.addEventListener('alarm', (e) => {
lastEventId = e.lastEventId;
const data = JSON.parse(e.data);
console.log('告警:', data);
});
es.addEventListener('error', () => {
// 浏览器会自动重连,并附带 Last-Event-ID 头
console.log('等待自动重连...');
});
}
connect();
cURL 示例
curl -N -H "Authorization: Bearer <token>" \
"http://localhost:8080/api/v1/events/stream?topics=alarm&severity=critical"
-N 参数禁用缓冲,确保事件实时输出。
Python 示例
import requests
import json
url = "http://localhost:8080/api/v1/events/stream"
headers = {"Authorization": "Bearer <token>"}
params = {"topics": "alarm,agent", "severity": "critical"}
response = requests.get(url, headers=headers, params=params, stream=True)
for line in response.iter_lines(decode_unicode=True):
if not line:
continue
if line.startswith("event:"):
event_type = line[6:].strip()
elif line.startswith("id:"):
event_id = line[3:].strip()
elif line.startswith("data:"):
data = json.loads(line[5:].strip())
print(f"[{event_type}] {data}")
Rust 示例
use reqwest::header;
use futures_util::StreamExt;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let client = reqwest::Client::new();
let mut headers = header::HeaderMap::new();
headers.insert(
header::AUTHORIZATION,
header::HeaderValue::from_static("Bearer <token>"),
);
let response = client
.get("http://localhost:8080/api/v1/events/stream")
.headers(headers)
.query(&[("topics", "alarm"), ("severity", "critical")])
.send()
.await?;
let mut stream = response.bytes_stream();
let mut buffer = String::new();
while let Some(chunk) = stream.next().await {
let chunk = chunk?;
buffer.push_str(&String::from_utf8_lossy(&chunk));
while let Some(pos) = buffer.find("\n\n") {
let event_str = buffer[..pos].to_string();
buffer = buffer[pos + 2..].to_string();
for line in event_str.lines() {
if line.starts_with("event:") {
println!("类型: {}", line[6..].trim());
} else if line.starts_with("data:") {
println!("数据: {}", line[5..].trim());
} else if line.starts_with("id:") {
println!("ID: {}", line[3..].trim());
}
}
println!("---");
}
}
Ok(())
}
JavaScript 完整看板示例
class EnerOSDashboard {
constructor(token) {
this.token = token;
this.es = null;
this.alarmCount = 0;
this.agentStatus = new Map();
}
start() {
const url = new URL('http://localhost:8080/api/v1/events/stream');
url.searchParams.set('topics', 'alarm,agent,topology,constraint');
url.searchParams.set('token', this.token);
this.es = new EventSource(url.toString());
this.es.addEventListener('alarm', (e) => this.onAlarm(e));
this.es.addEventListener('agent', (e) => this.onAgent(e));
this.es.addEventListener('topology', (e) => this.onTopology(e));
this.es.addEventListener('constraint', (e) => this.onConstraint(e));
this.es.addEventListener('system', (e) => this.onSystem(e));
this.es.addEventListener('open', () => {
console.log('[Dashboard] 已连接');
});
this.es.addEventListener('error', (e) => {
if (this.es.readyState === EventSource.CLOSED) {
console.error('[Dashboard] 连接已关闭,5s 后重连');
setTimeout(() => this.start(), 5000);
}
});
}
onAlarm(e) {
const data = JSON.parse(e.data);
this.alarmCount++;
console.log(`[告警 #${this.alarmCount}] ${data.severity}: ${data.message}`);
}
onAgent(e) {
const data = JSON.parse(e.data);
this.agentStatus.set(data.agent_id, data.new_status);
console.log(`[Agent] ${data.agent_id}: ${data.old_status} → ${data.new_status}`);
}
onTopology(e) {
const data = JSON.parse(e.data);
console.log(`[拓扑] ${data.network_id}:`, data.changes);
}
onConstraint(e) {
const data = JSON.parse(e.data);
console.log(`[约束] ${data.element} ${data.quantity}: ${data.value} / ${data.limit}`);
}
onSystem(e) {
const data = JSON.parse(e.data);
console.log(`[系统] ${data.type}: ${data.message}`);
}
stop() {
if (this.es) {
this.es.close();
this.es = null;
}
}
}
const dashboard = new EnerOSDashboard('<bearer-token>');
dashboard.start();
重连机制
浏览器自动重连
浏览器 EventSource 在连接断开后自动重连,默认间隔 3 秒。服务端可通过 retry 字段调整:
retry: 5000
Last-Event-ID 断点续传
浏览器在重连时会自动附带 Last-Event-ID HTTP 头,包含最后接收的事件 ID。服务端据此补发遗漏事件:
GET /api/v1/events/stream?topics=alarm HTTP/1.1
Last-Event-ID: evt_20260706_001
服务端处理逻辑:
- 解析
Last-Event-IDHeader(或 Query 参数last_event_id) - 从事件存储中查询该 ID 之后的所有事件
- 按 ID 顺序补发
- 补发完成后,继续推送实时事件
重连退避策略
第 1 次重连:3s
第 2 次重连:3s
第 3 次重连:6s
第 4 次重连:12s
第 5 次及以上:30s(上限)
服务端通过 retry 字段动态调整:
retry: 3000
event: system
data: {"type":"reconnect_backoff","interval_ms":6000}
SSE vs WebSocket 对比
| 维度 | SSE | WebSocket |
|---|---|---|
| 通信方向 | 单向(服务端→客户端) | 双向 |
| 协议 | HTTP/1.1 或 HTTP/2 | WS/WSS |
| 浏览器支持 | 原生 EventSource | 原生 WebSocket |
| 自动重连 | 是(浏览器内置) | 需手动实现 |
| 断点续传 | 原生支持(Last-Event-ID) | 需自定义协议 |
| 代理穿透 | 好(标准 HTTP) | 一般(需 Upgrade 头) |
| 二进制支持 | 否 | 是 |
| 最大连接数 | 浏览器单域名 6 个 | 无限制 |
| 协议开销 | 低(首包后仅文本) | 低(首包后二进制) |
| 心跳实现 | 注释行 :keepalive | ping/pong 帧 |
| 适用场景 | 告警/事件看板、通知推送 | Agent 指令交互、实时控制 |
何时选择 SSE
- 仅需服务端推送,无需客户端发送数据
- 需要浏览器原生自动重连
- 部署环境代理穿透要求高
- 告警看板、事件大屏、运维通知
何时选择 WebSocket
- 需要客户端向服务端发送指令
- 需要传输二进制数据
- 需要更高的实时性(毫秒级)
- Agent 远程控制、双向协作
限制
| 维度 | 限制值 |
|---|---|
| 单连接订阅主题 | 10 个 |
| 事件 payload 大小 | 64 KB |
| 单 token 并发连接 | 5 个 |
| 单租户并发连接 | 50 个 |
| 浏览器单域名连接 | 6 个(HTTP/1.1) |
| 事件保留时长 | 24 小时 |
| 连接空闲超时 | 5 分钟(无事件且无心跳响应) |
错误处理
连接错误
| 场景 | HTTP 状态 | 处理建议 |
|---|---|---|
| 未认证 | 401 | 检查 Token |
| 权限不足 | 403 | 联系管理员 |
| 主题数超限 | 400 | 减少订阅主题 |
| 服务不可用 | 503 | 等待后重连 |
流内错误
SSE 连接建立后,错误通过 system 事件推送:
event: system
id: evt_err_001
data: {"type":"error","code":"SUBSCRIPTION_EXPIRED","message":"订阅已过期,请重新连接"}
性能调优
- 合理使用
topics过滤:避免订阅全部主题导致带宽浪费 - 使用 HTTP/2:多路复用,避免连接数限制
- 客户端限流:高频事件场景下,客户端应做防抖或节流
- 压缩传输:启用 gzip 压缩,减少带宽占用
- CDN 缓存:历史事件查询可通过 REST 端点 + CDN 缓存
相关文档
- API 总览 — 四种 API 对比
- WebSocket — 双向实时通信
- 限流与配额 — 连接数限制
- REST API — 历史事件查询
- eneros-eventbus crate — 事件总线实现