跳到主内容

SSE 实时推送

API 参考

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>
字段类型必填说明
eventstring事件类型,客户端可按类型分发;缺省为 message
idstring事件 ID,浏览器存储后重连时通过 Last-Event-ID 头发送
retryint重连间隔(毫秒),浏览器将覆盖默认值
datastring事件数据,多行 data: 将以换行符拼接

: 开头的行为注释行,常用于心跳保活。

端点信息

项目取值
URLGET http://localhost:8080/api/v1/events/stream
Content-Typetext/event-stream; charset=utf-8
Cache-Controlno-cache
Connectionkeep-alive
认证Bearer Token(Header)或 Query 参数 ?token=
心跳周期15 秒(发送 :keepalive\n\n 注释行)
重连间隔默认 3 秒,可通过 retry 字段调整
单连接最大订阅10 个主题
单租户最大连接50 个并发
事件 payload 上限64 KB
事件保留时长24 小时(用于断点续传)

订阅参数

通过 Query 参数过滤事件流:

参数类型必填默认值说明
topicsstring全部逗号分隔的主题列表,如 alarm,agent
network_idstring限定电网 ID
severitystring告警级别过滤:info/warn/error/critical
agent_idstring限定 Agent ID
sincedatetime起始时间(RFC 3339)
untildatetime结束时间(达到后自动关闭连接)
last_event_idstring上次接收的事件 ID,用于断点续传
formatstringjsonjson / plain
include_auditboolfalse是否包含审计事件(需审计权限)

事件类型

SSE 端点推送以下事件类型,每个 event 字段对应一个分类:

event说明payload 关键字段触发频率
alarm告警产生/确认/清除alarm_id, severity, message实时
agentAgent 状态变更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心跳事件(应用层)timestamp60 秒

消息格式示例

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

服务端处理逻辑:

  1. 解析 Last-Event-ID Header(或 Query 参数 last_event_id
  2. 从事件存储中查询该 ID 之后的所有事件
  3. 按 ID 顺序补发
  4. 补发完成后,继续推送实时事件

重连退避策略

第 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 对比

维度SSEWebSocket
通信方向单向(服务端→客户端)双向
协议HTTP/1.1 或 HTTP/2WS/WSS
浏览器支持原生 EventSource原生 WebSocket
自动重连是(浏览器内置)需手动实现
断点续传原生支持(Last-Event-ID)需自定义协议
代理穿透好(标准 HTTP)一般(需 Upgrade 头)
二进制支持
最大连接数浏览器单域名 6 个无限制
协议开销低(首包后仅文本)低(首包后二进制)
心跳实现注释行 :keepaliveping/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 缓存

相关文档