说说【Flask SSE 实时大屏】从数据库新增记录到 ECharts 更新的完整演示链路(学习笔记)

这两天一直在研究这个话题,踩了几个坑,把遇到的东西整理成文,供有需要的朋友参考。

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

🏠博主简介

文章目录

前言

监控大屏在原型阶段常用 Math.random() 验证图表更新。进入数据链路演示时,随机数应从页面移除:采集接口负责写入,数据库保存记录,Flask 用 SSE 推送新增数据,ECharts 只展示收到的内容。

下面给出一套可以直接运行和逐层检查的示例,并补充断线续传、代理缓冲、连接扩展以及访问控制边界。文中不声称经过生产压测。

示例使用 Flask 和 SQLite,方便直接运行。换成 MySQL、PostgreSQL 或真实消息队列后,分层思路不变,但查询和并发方案需要按实际环境调整。


一、先删掉“假实时”的数据源

原页面每两秒执行一次:

setInterval(() => {
  const value = 20 + Math.random() * 10;
  appendPoint(new Date(), value);
}, 2000);

这段代码用来确认图表能否更新没有问题,但它不应该继续留在正式数据链路里。否则会出现几个麻烦:

  • 页面有数据,数据库却查不到对应记录;
  • 前端刷新后,历史曲线全部消失;
  • 无法判断异常值来自设备还是随机数;
  • 多个浏览器看到的曲线互不相同;
  • 后端接口即使故障,页面仍然显得“正常”。
页面数据只保留两个入口:

1. 首次打开时,通过历史接口加载最近一段记录;
2. 页面打开后,通过 SSE 接收新增记录。

测试数据也一定要走采集接口写入,不能直接塞进 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  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}'

验证时按下面的顺序检查:

1. POST 返回的 id 是否递增;
2. SQLite 中能否查到同一条记录;
3. 浏览器网络面板中的 /stream 是否保持连接;
4. 事件流里是否出现对应 id 和温度;
5. ECharts 曲线是否只新增一个点;
6. 刷新页面后,历史接口是否能恢复刚才的数据。

这样一条记录从入口到页面都有证据,不需要靠“图动了”判断系统是否正常。


九、部署到 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”的发现机制,并补齐鉴权、租户过滤、跨域策略、连接容量与可观测性。


参考资料


今天的内容大概就这些,实际开发中大家还会遇到更多细节,欢迎留言分享自己的经验。

评论 (0)

暂无评论