
🔥个人主页:爱和冰阔乐 📚专栏传送门:《数据结构与算法》 、C++ 🐶学习方向:C++方向学习爱好者 ⭐人生格言:得知坦然 ,失之淡然

🏠博主简介 
文章目录
-
- 前言
- 一、先删掉“假实时”的数据源
- 二、为什么这个演示选择 SSE
- 三、先把数据库表设计成事实来源
- 四、采集接口只负责校验和落库
- 五、补一个历史接口,避免页面刚打开时是空白
- 六、Flask 推送端的完整写法
- 七、ECharts 只接收数据,不再自己编数据
- 八、不写模拟器,也能验证整条链路
- 九、部署到 Nginx 后,消息为什么突然成批出现
- 十、这个简单版本不能直接承担高并发
- 十一、鉴权、租户隔离与 CORS
- 十二、四层职责与排查入口
- 链路边界
- 参考资料
前言
监控大屏在原型阶段常用 Math.random() 验证图表更新。进入数据链路演示时,随机数应从页面移除:采集接口负责写入,数据库保存记录,Flask 用 SSE 推送新增数据,ECharts 只展示收到的内容。
下面给出一套可以直接运行和逐层检查的示例,并补充断线续传、代理缓冲、连接扩展以及访问控制边界。文中不声称经过生产压测。
示例使用 Flask 和 SQLite,方便直接运行。换成 MySQL、PostgreSQL 或真实消息队列后,分层思路不变,但查询和并发方案需要按实际环境调整。
一、先删掉“假实时”的数据源
原页面每两秒执行一次:
setInterval(() => {
const value = 20 + Math.random() * 10;
appendPoint(new Date(), value);
}, 2000);
这段代码用来确认图表能否更新没有问题,但它不应该继续留在正式数据链路里。否则会出现几个麻烦:
- 页面有数据,数据库却查不到对应记录;
- 前端刷新后,历史曲线全部消失;
- 无法判断异常值来自设备还是随机数;
- 多个浏览器看到的曲线互不相同;
- 后端接口即使故障,页面仍然显得“正常”。
页面数据只保留两个入口:
测试数据也必须走采集接口写入,不能直接塞进 ECharts。
删除随机数以后,一条数据会依次经过采集、存储、推送和展示。

二、为什么这个演示选择 SSE
可以先比较三种方式:
| 定时轮询 | 实现最简单,但无变化时也会重复请求 | 可以用,但请求较多 |
| WebSocket | 双向通信,适合频繁交互 | 能用,但这次不需要浏览器向服务端持续发消息 |
| SSE | 服务端单向推送,浏览器原生支持事件流 | 符合当前数据方向 |
这个页面只需要服务器把新遥测记录推给浏览器,没有聊天、控制指令或二进制数据,因此 SSE 足以表达当前通信方向。
SSE 的响应类型是 text/event-stream,每条消息用空行分隔。浏览器可以用 EventSource 建立连接,并在连接中断后自动尝试重连。
它不是比 WebSocket 更“高级”的方案,只是这次的通信方向更简单。技术选择要看通信方向,而不是看哪个名称更“实时”。以后如果页面增加实时控制、双向状态同步,再重新评估 WebSocket 更合适。
三、先把数据库表设计成事实来源
示例表只保留设备编号、温度、湿度和采集时间:
CREATE TABLE IF NOT EXISTS telemetry (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device_id TEXT NOT NULL,
temperature REAL NOT NULL,
humidity REAL NOT NULL,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_telemetry_device_id
ON telemetry (device_id, id);
这里使用递增 id 作为推送游标。相比只使用时间,它更容易判断“上次已经发送到哪一条”,也不会因为两条记录采集时间相同而漏掉数据。
时间仍然由采集请求携带或由后端生成,用于图表横轴和业务分析;id 主要服务于数据库顺序和断线续传。
初始化代码如下:
from pathlib import Path
import sqlite3
DB_PATH = Path(__file__).with_name("telemetry.db")
def init_db():
with sqlite3.connect(DB_PATH) as conn:
conn.executescript(
"""
CREATE TABLE IF NOT EXISTS telemetry (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device_id TEXT NOT NULL,
temperature REAL NOT NULL,
humidity REAL NOT NULL,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_telemetry_device_id
ON telemetry (device_id, id);
"""
)
四、采集接口只负责校验和落库
设备或采集程序向下面的接口提交数据:
POST /api/telemetry
Flask 代码:
from datetime import datetime, timezone
import sqlite3
from flask import Flask, jsonify, request
app = Flask(__name__)
@app.post("/api/telemetry")
def create_telemetry():
body = request.get_json(silent=True) or {}
device_id = str(body.get("device_id", "")).strip()
if not device_id:
return jsonify({"message": "device_id 不能为空"}), 400
try:
temperature = float(body["temperature"])
humidity = float(body["humidity"])
except (KeyError, TypeError, ValueError):
return jsonify({"message": "温度或湿度格式错误"}), 400
if not –50 <= temperature <= 100:
return jsonify({"message": "温度超出允许范围"}), 400
if not 0 <= humidity <= 100:
return jsonify({"message": "湿度超出允许范围"}), 400
created_at = body.get("created_at")
if created_at is None:
created_at = datetime.now(timezone.utc).isoformat()
with sqlite3.connect(DB_PATH) as conn:
cursor = conn.execute(
"""
INSERT INTO telemetry (
device_id, temperature, humidity, created_at
) VALUES (?, ?, ?, ?)
""",
(device_id, temperature, humidity, created_at),
)
row_id = cursor.lastrowid
return jsonify({"id": row_id}), 201
示例为了突出链路,只对时间做了最简单处理。接入设备时还要规定时间格式、时区和可信程度。如果设备时钟可能漂移,可以同时保存“设备采集时间”和“服务器接收时间”,避免让一个字段承担两种含义。
接口返回 201 后,记录已经成为数据库里的事实。SSE 推送端只读取已落库数据,不直接消费某个全局 Python 变量。这样即使页面晚一点打开,也可以从历史记录恢复。
五、补一个历史接口,避免页面刚打开时是空白
SSE 适合传递新增记录,但页面首次打开通常还需要最近几十个点。
@app.get("/api/telemetry/history")
def telemetry_history():
device_id = request.args.get("device_id", "sensor-01").strip()
try:
limit = min(max(int(request.args.get("limit", 60)), 1), 500)
except ValueError:
return jsonify({"message": "limit 格式错误"}), 400
with sqlite3.connect(DB_PATH) as conn:
conn.row_factory = sqlite3.Row
rows = conn.execute(
"""
SELECT id, device_id, temperature, humidity, created_at
FROM telemetry
WHERE device_id = ?
ORDER BY id DESC
LIMIT ?
""",
(device_id, limit),
).fetchall()
items = [dict(row) for row in reversed(rows)]
return jsonify({"items": items})
查询先倒序取最近 limit 条,再在返回前恢复成正序。ECharts 得到的数据因此会从旧到新排列。
limit 在后端设置上限,避免有人传入一个很大的数字,一次把整张表都取走。
历史数据和实时事件需要在同一个 ID 位置完成衔接。

六、Flask 推送端的完整写法
SSE 的消息格式并不复杂:
id: 128
event: telemetry
data: {"temperature": 25.6}
最后的空行不能省略。完整接口如下:
import json
import time
from flask import Response, request, stream_with_context
@app.get("/stream")
def telemetry_stream():
device_id = request.args.get("device_id", "sensor-01").strip()
last_event_id = request.headers.get("Last-Event-ID")
fallback_id = request.args.get("after_id", "0")
try:
start_id = int(last_event_id or fallback_id)
except ValueError:
start_id = 0
@stream_with_context
def generate():
cursor_id = start_id
while True:
with sqlite3.connect(DB_PATH) as conn:
conn.row_factory = sqlite3.Row
rows = conn.execute(
"""
SELECT id, device_id, temperature, humidity, created_at
FROM telemetry
WHERE device_id = ? AND id > ?
ORDER BY id
LIMIT 100
""",
(device_id, cursor_id),
).fetchall()
if rows:
for row in rows:
item = dict(row)
cursor_id = item["id"]
data = json.dumps(item, ensure_ascii=False)
yield (
f"id: {cursor_id}\\n"
"event: telemetry\\n"
f"data: {data}\\n\\n"
)
else:
# 注释行可作为心跳,避免中间代理长时间看不到数据
yield ": keep-alive\\n\\n"
time.sleep(1)
return Response(
generate(),
mimetype="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
},
)
这段代码有几个刻意保留的细节。
第一,每轮查询都使用独立、短生命周期的数据库连接,没有把请求外部创建的游标一直挂在长连接上。
第二,每条事件都发送 id。浏览器重连时可以携带 Last-Event-ID,后端据此继续读取后续记录。
第三,没有新数据时发送注释形式的心跳。SSE 客户端会忽略以冒号开头的行,但它能让代理知道连接仍然活着。
第四,响应禁止缓存,并通过 X-Accel-Buffering: no 提示 Nginx 不要攒够一批再返回。不过正式部署时仍应在代理配置中显式关闭缓冲,不能只依赖这个响应头。
当连接中断时,事件 ID 还承担续传位置的作用。

每条事件发送稳定 ID,重连时才能从已确认位置继续,而不是重新推送全部历史。
七、ECharts 只接收数据,不再自己编数据
页面初始化:
<div id="temperature-chart" style="height: 360px"></div>
<div id="stream-status">正在连接</div>
const chart = echarts.init(
document.getElementById('temperature-chart')
);
const statusElement = document.getElementById('stream-status');
const points = [];
chart.setOption({
title: { text: '设备温度' },
tooltip: { trigger: 'axis' },
xAxis: { type: 'time' },
yAxis: { type: 'value', name: '℃' },
series: [
{
name: '温度',
type: 'line',
showSymbol: false,
data: points
}
]
});
先加载历史数据:
async function loadHistory() {
const response = await fetch(
'/api/telemetry/history?device_id=sensor-01&limit=60'
);
if (!response.ok) {
throw new Error(`历史数据加载失败:${response.status}`);
}
const payload = await response.json();
points.length = 0;
for (const item of payload.items) {
points.push([item.created_at, item.temperature]);
}
chart.setOption({
series: [{ data: points }]
});
return payload.items.at(–1)?.id ?? 0;
}
再从历史数据最后一条开始订阅:
async function startStream() {
const lastId = await loadHistory();
const source = new EventSource(
`/stream?device_id=sensor-01&after_id=${lastId}`
);
source.onopen = () => {
statusElement.textContent = '实时连接正常';
};
source.addEventListener('telemetry', (event) => {
const item = JSON.parse(event.data);
points.push([item.created_at, item.temperature]);
if (points.length > 60) {
points.shift();
}
chart.setOption({
series: [{ data: points }]
});
});
source.onerror = () => {
statusElement.textContent = '连接中断,正在重试';
};
window.addEventListener('beforeunload', () => {
source.close();
});
}
startStream().catch((error) => {
statusElement.textContent = error.message;
});
如果图表同时展示很多系列,不建议每来一个点就重建整个 option。可以只更新发生变化的 series.data,并控制内存中的窗口长度。ECharts 的 setOption 会按配置更新图表,不需要销毁后重新初始化。
八、不写模拟器,也能验证整条链路
为了避免测试代码混入页面,可以用 curl 向采集接口写入几条确定数据。
curl -X POST http://127.0.0.1:5000/api/telemetry \\
-H 'Content-Type: application/json' \\
-d '{"device_id":"sensor-01","temperature":24.8,"humidity":51.2}'
再写一条:
curl -X POST http://127.0.0.1:5000/api/telemetry \\
-H 'Content-Type: application/json' \\
-d '{"device_id":"sensor-01","temperature":25.1,"humidity":50.7}'
验证时按下面的顺序检查:
这样一条记录从入口到页面都有证据,不需要靠“图动了”判断系统是否正常。
九、部署到 Nginx 后,消息为什么突然成批出现
如果直连 Flask 时事件逐条到达,经过 Nginx 后却成批出现,应优先检查代理是否缓冲上游响应,再检查前端图表。
关闭缓冲和继续缓冲时,浏览器收到事件的节奏会不同。

对应位置需要关闭缓冲并延长读取超时:
location /stream {
proxy_pass http://app:5000;
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_buffering off;
proxy_cache off;
proxy_read_timeout 1h;
}
修改后先执行配置检查,再平滑重载:
nginx -t
nginx -s reload
如果前面还有 CDN、网关或其他反向代理,也要逐层确认它们是否支持并保留流式响应。只改最靠近 Flask 的一层,不代表整条链路都不会缓冲。
十、这个简单版本不能直接承担高并发
示例为了容易理解,每个 SSE 客户端每秒查询一次 SQLite。每增加一个页面,就会增加一条长期连接和一个重复查询循环。这种“每客户端每秒查询 SQLite”的演示结构不适合横向扩展,也不能当作生产架构或压测结论。
进入实际部署前至少要评估下面几项:
- WSGI 服务和 worker 类型能否承受长期连接;
- 每个客户端轮询数据库是否造成重复读取;
- 多进程实例之间怎样共享新增事件;
- 客户端断线期间积累的数据怎样补发;
- 是否需要按用户、租户或设备做权限校验;
- 连接数、发送延迟和断线率怎样监控。
数据量和连接数增加后,可以把“发现新增记录”的职责移给消息系统,例如 Redis Streams 或专门的消息队列。数据库仍保存业务事实,不再由每个浏览器连接单独查询同一张表。
另外,SSE 是长期 HTTP 连接。浏览器对同一域名的连接数限制、HTTP 版本和页面打开数量都可能影响表现。如果一个页面为每张图单独建立一个 SSE 连接,很快就会浪费连接资源。更合理的方式是一个页面共用一条事件流,再在前端按事件类型分发。
十一、鉴权、租户隔离与 CORS
SSE 使用长期 HTTP 请求,不会自动绕过鉴权。历史接口和 /stream 必须执行同一套权限判断,device_id 也不能由客户端任意指定后直接查询。
多租户场景中,tenant_id 应从已经验证的会话或令牌中取得,并同时加入历史与实时查询条件:
SELECT id, device_id, temperature, humidity, created_at
FROM telemetry
WHERE tenant_id = ?
AND device_id = ?
AND id > ?
ORDER BY id
LIMIT 100;
数据库索引也应与访问路径匹配,例如 (tenant_id, device_id, id)。只校验“设备存在”还不够,还要证明该设备属于当前租户。
原生 EventSource 没有通用的自定义请求头入口。同源部署可以使用受保护的会话 Cookie;跨源并需要 Cookie 时,客户端可以显式启用凭据:
const source = new EventSource(
'https://api.example.com/stream?device_id=sensor-01',
{ withCredentials: true }
);
服务端必须返回明确的允许来源和凭据头;携带凭据时不能把 Access-Control-Allow-Origin 设置为 *。CORS 只决定浏览器能否读取跨源响应,不能替代身份验证和租户授权。
也不建议把长期有效的访问令牌直接放在查询参数中,因为 URL 可能进入浏览器历史、代理日志和监控系统。必须跨域时,应结合部署条件选择短期票据、同源反向代理或能够安全携带凭据的客户端方案。
十二、四层职责与排查入口
完成改造后,四层职责变得很清楚:
| 采集接口 | 校验设备数据并写入数据库 |
| 数据库 | 保存可以追溯的历史事实 |
| SSE 接口 | 按递增 ID 推送新增记录并支持重连续传 |
| ECharts 页面 | 加载历史、接收新增数据并更新视图 |
任何一层出问题,都可以单独检查。
页面没变化时,先看采集接口有没有写入;数据库有记录但事件流没有,就检查推送查询和代理;事件流有消息但图表没动,再看 JSON 解析和 ECharts 更新。
把四层各自能观察到的证据列开以后,故障位置会清楚很多。

相比把所有逻辑塞进一个定时器,这条链路更长,却更容易定位问题。
链路边界
演示页面的数据先通过采集接口落库,首次打开时加载历史,随后接收新增事件。每条 SSE 事件携带 ID,断线重连时可以续传;代理层关闭缓冲,ECharts 只维护有限窗口。
这套代码用于说明协议和排查方法。扩展到更多用户和设备时,需要替换“每客户端轮询 SQLite”的发现机制,并补齐鉴权、租户过滤、跨域策略、连接容量与可观测性。
参考资料
- Flask Documentation:Streaming Contents
- MDN Web Docs:Using server-sent events
- MDN Web Docs:CORS credentials and wildcard origins
- Apache ECharts Handbook:Dynamic Data



