个人主页爱和冰阔乐专栏传送门《数据结构与算法》 、C学习方向C方向学习爱好者⭐人生格言得知坦然 失之淡然博主简介文章目录前言一、先删掉“假实时”的数据源二、为什么这个演示选择 SSE三、先把数据库表设计成事实来源四、采集接口只负责校验和落库五、补一个历史接口避免页面刚打开时是空白六、Flask 推送端的完整写法七、ECharts 只接收数据不再自己编数据八、不写模拟器也能验证整条链路九、部署到 Nginx 后消息为什么突然成批出现十、这个简单版本不能直接承担高并发十一、鉴权、租户隔离与 CORS十二、四层职责与排查入口链路边界参考资料前言监控大屏在原型阶段常用Math.random()验证图表更新。进入数据链路演示时随机数应从页面移除采集接口负责写入数据库保存记录Flask 用 SSE 推送新增数据ECharts 只展示收到的内容。下面给出一套可以直接运行和逐层检查的示例并补充断线续传、代理缓冲、连接扩展以及访问控制边界。文中不声称经过生产压测。示例使用 Flask 和 SQLite方便直接运行。换成 MySQL、PostgreSQL 或真实消息队列后分层思路不变但查询和并发方案需要按实际环境调整。一、先删掉“假实时”的数据源原页面每两秒执行一次setInterval((){constvalue20Math.random()*10;appendPoint(newDate(),value);},2000);这段代码用来确认图表能否更新没有问题但它不应该继续留在正式数据链路里。否则会出现几个麻烦页面有数据数据库却查不到对应记录前端刷新后历史曲线全部消失无法判断异常值来自设备还是随机数多个浏览器看到的曲线互不相同后端接口即使故障页面仍然显得“正常”。页面数据只保留两个入口首次打开时通过历史接口加载最近一段记录页面打开后通过 SSE 接收新增记录。测试数据也必须走采集接口写入不能直接塞进 ECharts。删除随机数以后一条数据会依次经过采集、存储、推送和展示。二、为什么这个演示选择 SSE可以先比较三种方式方式特点是否适合这次定时轮询实现最简单但无变化时也会重复请求可以用但请求较多WebSocket双向通信适合频繁交互能用但这次不需要浏览器向服务端持续发消息SSE服务端单向推送浏览器原生支持事件流符合当前数据方向这个页面只需要服务器把新遥测记录推给浏览器没有聊天、控制指令或二进制数据因此 SSE 足以表达当前通信方向。SSE 的响应类型是text/event-stream每条消息用空行分隔。浏览器可以用EventSource建立连接并在连接中断后自动尝试重连。它不是比 WebSocket 更“高级”的方案只是这次的通信方向更简单。技术选择要看通信方向而不是看哪个名称更“实时”。以后如果页面增加实时控制、双向状态同步再重新评估 WebSocket 更合适。三、先把数据库表设计成事实来源示例表只保留设备编号、温度、湿度和采集时间CREATETABLEIFNOTEXISTStelemetry(idINTEGERPRIMARYKEYAUTOINCREMENT,device_idTEXTNOTNULL,temperatureREALNOTNULL,humidityREALNOTNULL,created_atTEXTNOTNULL);CREATEINDEXIFNOTEXISTSidx_telemetry_device_idONtelemetry(device_id,id);这里使用递增id作为推送游标。相比只使用时间它更容易判断“上次已经发送到哪一条”也不会因为两条记录采集时间相同而漏掉数据。时间仍然由采集请求携带或由后端生成用于图表横轴和业务分析id主要服务于数据库顺序和断线续传。初始化代码如下frompathlibimportPathimportsqlite3 DB_PATHPath(__file__).with_name(telemetry.db)definit_db():withsqlite3.connect(DB_PATH)asconn: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/telemetryFlask 代码fromdatetimeimportdatetime,timezoneimportsqlite3fromflaskimportFlask,jsonify,request appFlask(__name__)app.post(/api/telemetry)defcreate_telemetry():bodyrequest.get_json(silentTrue)or{}device_idstr(body.get(device_id,)).strip()ifnotdevice_id:returnjsonify({message:device_id 不能为空}),400try:temperaturefloat(body[temperature])humidityfloat(body[humidity])except(KeyError,TypeError,ValueError):returnjsonify({message:温度或湿度格式错误}),400ifnot-50temperature100:returnjsonify({message:温度超出允许范围}),400ifnot0humidity100:returnjsonify({message:湿度超出允许范围}),400created_atbody.get(created_at)ifcreated_atisNone:created_atdatetime.now(timezone.utc).isoformat()withsqlite3.connect(DB_PATH)asconn:cursorconn.execute( INSERT INTO telemetry ( device_id, temperature, humidity, created_at ) VALUES (?, ?, ?, ?) ,(device_id,temperature,humidity,created_at),)row_idcursor.lastrowidreturnjsonify({id:row_id}),201示例为了突出链路只对时间做了最简单处理。接入设备时还要规定时间格式、时区和可信程度。如果设备时钟可能漂移可以同时保存“设备采集时间”和“服务器接收时间”避免让一个字段承担两种含义。接口返回201后记录已经成为数据库里的事实。SSE 推送端只读取已落库数据不直接消费某个全局 Python 变量。这样即使页面晚一点打开也可以从历史记录恢复。五、补一个历史接口避免页面刚打开时是空白SSE 适合传递新增记录但页面首次打开通常还需要最近几十个点。app.get(/api/telemetry/history)deftelemetry_history():device_idrequest.args.get(device_id,sensor-01).strip()try:limitmin(max(int(request.args.get(limit,60)),1),500)exceptValueError:returnjsonify({message:limit 格式错误}),400withsqlite3.connect(DB_PATH)asconn:conn.row_factorysqlite3.Row rowsconn.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)forrowinreversed(rows)]returnjsonify({items:items})查询先倒序取最近limit条再在返回前恢复成正序。ECharts 得到的数据因此会从旧到新排列。limit在后端设置上限避免有人传入一个很大的数字一次把整张表都取走。历史数据和实时事件需要在同一个 ID 位置完成衔接。六、Flask 推送端的完整写法SSE 的消息格式并不复杂id: 128 event: telemetry data: {temperature: 25.6}最后的空行不能省略。完整接口如下importjsonimporttimefromflaskimportResponse,request,stream_with_contextapp.get(/stream)deftelemetry_stream():device_idrequest.args.get(device_id,sensor-01).strip()last_event_idrequest.headers.get(Last-Event-ID)fallback_idrequest.args.get(after_id,0)try:start_idint(last_event_idorfallback_id)exceptValueError:start_id0stream_with_contextdefgenerate():cursor_idstart_idwhileTrue:withsqlite3.connect(DB_PATH)asconn:conn.row_factorysqlite3.Row rowsconn.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()ifrows:forrowinrows:itemdict(row)cursor_iditem[id]datajson.dumps(item,ensure_asciiFalse)yield(fid:{cursor_id}\nevent: telemetry\nfdata:{data}\n\n)else:# 注释行可作为心跳避免中间代理长时间看不到数据yield: keep-alive\n\ntime.sleep(1)returnResponse(generate(),mimetypetext/event-stream,headers{Cache-Control:no-cache,X-Accel-Buffering:no,},)这段代码有几个刻意保留的细节。第一每轮查询都使用独立、短生命周期的数据库连接没有把请求外部创建的游标一直挂在长连接上。第二每条事件都发送id。浏览器重连时可以携带Last-Event-ID后端据此继续读取后续记录。第三没有新数据时发送注释形式的心跳。SSE 客户端会忽略以冒号开头的行但它能让代理知道连接仍然活着。第四响应禁止缓存并通过X-Accel-Buffering: no提示 Nginx 不要攒够一批再返回。不过正式部署时仍应在代理配置中显式关闭缓冲不能只依赖这个响应头。当连接中断时事件 ID 还承担续传位置的作用。每条事件发送稳定 ID重连时才能从已确认位置继续而不是重新推送全部历史。七、ECharts 只接收数据不再自己编数据页面初始化dividtemperature-chartstyleheight:360px/divdividstream-status正在连接/divconstchartecharts.init(document.getElementById(temperature-chart));conststatusElementdocument.getElementById(stream-status);constpoints[];chart.setOption({title:{text:设备温度},tooltip:{trigger:axis},xAxis:{type:time},yAxis:{type:value,name:℃},series:[{name:温度,type:line,showSymbol:false,data:points}]});先加载历史数据asyncfunctionloadHistory(){constresponseawaitfetch(/api/telemetry/history?device_idsensor-01limit60);if(!response.ok){thrownewError(历史数据加载失败${response.status});}constpayloadawaitresponse.json();points.length0;for(constitemofpayload.items){points.push([item.created_at,item.temperature]);}chart.setOption({series:[{data:points}]});returnpayload.items.at(-1)?.id??0;}再从历史数据最后一条开始订阅asyncfunctionstartStream(){constlastIdawaitloadHistory();constsourcenewEventSource(/stream?device_idsensor-01after_id${lastId});source.onopen(){statusElement.textContent实时连接正常;};source.addEventListener(telemetry,(event){constitemJSON.parse(event.data);points.push([item.created_at,item.temperature]);if(points.length60){points.shift();}chart.setOption({series:[{data:points}]});});source.onerror(){statusElement.textContent连接中断正在重试;};window.addEventListener(beforeunload,(){source.close();});}startStream().catch((error){statusElement.textContenterror.message;});如果图表同时展示很多系列不建议每来一个点就重建整个 option。可以只更新发生变化的series.data并控制内存中的窗口长度。ECharts 的setOption会按配置更新图表不需要销毁后重新初始化。八、不写模拟器也能验证整条链路为了避免测试代码混入页面可以用curl向采集接口写入几条确定数据。curl-XPOST http://127.0.0.1:5000/api/telemetry\-HContent-Type: application/json\-d{device_id:sensor-01,temperature:24.8,humidity:51.2}再写一条curl-XPOST http://127.0.0.1:5000/api/telemetry\-HContent-Type: application/json\-d{device_id:sensor-01,temperature:25.1,humidity:50.7}验证时按下面的顺序检查POST 返回的id是否递增SQLite 中能否查到同一条记录浏览器网络面板中的/stream是否保持连接事件流里是否出现对应id和温度ECharts 曲线是否只新增一个点刷新页面后历史接口是否能恢复刚才的数据。这样一条记录从入口到页面都有证据不需要靠“图动了”判断系统是否正常。九、部署到 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-tnginx-sreload如果前面还有 CDN、网关或其他反向代理也要逐层确认它们是否支持并保留流式响应。只改最靠近 Flask 的一层不代表整条链路都不会缓冲。十、这个简单版本不能直接承担高并发示例为了容易理解每个 SSE 客户端每秒查询一次 SQLite。每增加一个页面就会增加一条长期连接和一个重复查询循环。这种“每客户端每秒查询 SQLite”的演示结构不适合横向扩展也不能当作生产架构或压测结论。进入实际部署前至少要评估下面几项WSGI 服务和 worker 类型能否承受长期连接每个客户端轮询数据库是否造成重复读取多进程实例之间怎样共享新增事件客户端断线期间积累的数据怎样补发是否需要按用户、租户或设备做权限校验连接数、发送延迟和断线率怎样监控。数据量和连接数增加后可以把“发现新增记录”的职责移给消息系统例如 Redis Streams 或专门的消息队列。数据库仍保存业务事实不再由每个浏览器连接单独查询同一张表。另外SSE 是长期 HTTP 连接。浏览器对同一域名的连接数限制、HTTP 版本和页面打开数量都可能影响表现。如果一个页面为每张图单独建立一个 SSE 连接很快就会浪费连接资源。更合理的方式是一个页面共用一条事件流再在前端按事件类型分发。十一、鉴权、租户隔离与 CORSSSE 使用长期 HTTP 请求不会自动绕过鉴权。历史接口和/stream必须执行同一套权限判断device_id也不能由客户端任意指定后直接查询。多租户场景中tenant_id应从已经验证的会话或令牌中取得并同时加入历史与实时查询条件SELECTid,device_id,temperature,humidity,created_atFROMtelemetryWHEREtenant_id?ANDdevice_id?ANDid?ORDERBYidLIMIT100;数据库索引也应与访问路径匹配例如(tenant_id, device_id, id)。只校验“设备存在”还不够还要证明该设备属于当前租户。原生EventSource没有通用的自定义请求头入口。同源部署可以使用受保护的会话 Cookie跨源并需要 Cookie 时客户端可以显式启用凭据constsourcenewEventSource(https://api.example.com/stream?device_idsensor-01,{withCredentials:true});服务端必须返回明确的允许来源和凭据头携带凭据时不能把Access-Control-Allow-Origin设置为*。CORS 只决定浏览器能否读取跨源响应不能替代身份验证和租户授权。也不建议把长期有效的访问令牌直接放在查询参数中因为 URL 可能进入浏览器历史、代理日志和监控系统。必须跨域时应结合部署条件选择短期票据、同源反向代理或能够安全携带凭据的客户端方案。十二、四层职责与排查入口完成改造后四层职责变得很清楚层职责采集接口校验设备数据并写入数据库数据库保存可以追溯的历史事实SSE 接口按递增 ID 推送新增记录并支持重连续传ECharts 页面加载历史、接收新增数据并更新视图任何一层出问题都可以单独检查。页面没变化时先看采集接口有没有写入数据库有记录但事件流没有就检查推送查询和代理事件流有消息但图表没动再看 JSON 解析和 ECharts 更新。把四层各自能观察到的证据列开以后故障位置会清楚很多。相比把所有逻辑塞进一个定时器这条链路更长却更容易定位问题。链路边界演示页面的数据先通过采集接口落库首次打开时加载历史随后接收新增事件。每条 SSE 事件携带 ID断线重连时可以续传代理层关闭缓冲ECharts 只维护有限窗口。这套代码用于说明协议和排查方法。扩展到更多用户和设备时需要替换“每客户端轮询 SQLite”的发现机制并补齐鉴权、租户过滤、跨域策略、连接容量与可观测性。参考资料Flask DocumentationStreaming ContentsMDN Web DocsUsing server-sent eventsMDN Web DocsCORS credentials and wildcard originsApache ECharts HandbookDynamic Data