欢迎光临
我们一直在努力

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

封面

🔥个人主页:爱和冰阔乐 📚专栏传送门:《数据结构与算法》 、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);

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

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

页面数据只保留两个入口:

  • 首次打开时,通过历史接口加载最近一段记录;
  • 页面打开后,通过 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 <= 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 位置完成衔接。

    历史接口与 SSE 的 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 还承担续传位置的作用。

    Last-Event-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}'

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

  • POST 返回的 id 是否递增;
  • SQLite 中能否查到同一条记录;
  • 浏览器网络面板中的 /stream 是否保持连接;
  • 事件流里是否出现对应 id 和温度;
  • ECharts 曲线是否只新增一个点;
  • 刷新页面后,历史接口是否能恢复刚才的数据。
  • 这样一条记录从入口到页面都有证据,不需要靠“图动了”判断系统是否正常。


    九、部署到 Nginx 后,消息为什么突然成批出现

    如果直连 Flask 时事件逐条到达,经过 Nginx 后却成批出现,应优先检查代理是否缓冲上游响应,再检查前端图表。

    关闭缓冲和继续缓冲时,浏览器收到事件的节奏会不同。

    Nginx 缓冲对 SSE 的影响

    对应位置需要关闭缓冲并延长读取超时:

    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
    赞(0)
    未经允许不得转载:171主机测评 » 【Flask SSE 实时大屏】从数据库新增记录到 ECharts 更新的完整演示链路
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址