ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

基于MQTT与SpringBoot构建实验室设备数据采集与实时监控系统

基于MQTT与SpringBoot构建实验室设备数据采集与实时监控系统 1. 项目概述与核心价值最近在实验室里折腾设备数据采集发现很多老设备还在用串口、Modbus这类传统协议数据孤岛现象严重想实时看个温湿度曲线都得手动记录效率太低。正好手头有个项目需要把十几台不同品牌、不同接口的仪器数据统一起来做个Web端的实时监控大屏。琢磨了一圈最终决定用MQTT协议做数据传输SpringBoot做后端服务搭建一套轻量级的实验室设备数据采集与实时监控系统。这套方案的核心思路就是把每台设备都变成一个独立的“发布者”通过MQTT这个高效的“消息快递员”把数据实时、异步地推送到云端服务端再由SpringBoot应用统一处理、存储并推送给前端大屏展示。为什么选这个组合首先实验室设备往往分散布线麻烦很多新设备支持网络但协议五花八门。MQTT协议基于发布/订阅模式特别适合物联网这种网络状况不稳定、设备性能参差不齐的场景。它协议头小功耗低一个TCP连接就能搞定双向通信比传统的HTTP轮询不停地问设备“你有新数据吗”省资源太多了。其次SpringBoot大家都熟快速构建Web服务的利器生态丰富整合MQTT客户端、数据库、WebSocket推流都非常方便。这个实战项目要解决的就是从零开始打通“设备端数据采集 - MQTT网络传输 - SpringBoot服务端汇聚与处理 - Web前端实时可视化”的全链路。这套系统适合谁呢如果你是实验室管理员、科研项目组的工程师或者是对物联网、实时系统感兴趣的后端/全栈开发者这个实战过程会很有参考价值。它不局限于特定设备核心是提供一套可复用的架构模式你完全可以根据自己的设备协议比如Modbus TCP、HTTP API、串口数据来适配数据采集端然后复用我们搭建的MQTTSpringBoot核心管道。2. 系统整体架构与核心组件选型2.1 架构设计思路拆解整个系统的设计遵循“解耦”和“实时”两个核心原则。我们不能让设备数据直接怼到数据库或者前端那样耦合太紧一个设备出问题可能拖垮整个系统。因此引入了MQTT消息中间件作为“数据总线”所有数据都通过它来流转。架构上分为四层设备与数据采集层这是数据的源头。可能是直接运行MQTT客户端代码的单片机如ESP32也可能是在工控机、树莓派上运行的数据采集程序负责从传感器、PLC或仪器仪表通过RS485/232、GPIO、网络API等读取数据并封装成JSON格式的消息发布到指定的MQTT主题。消息传输层由MQTT Broker服务器构成是整个系统的中枢神经。我们选用EMQX作为Broker它开源、高性能支持海量连接和集群社区活跃管理界面友好。它负责接收设备发布的消息并根据主题Topic将其转发给所有订阅了该主题的客户端。业务处理与存储层基于SpringBoot构建的核心后端服务。它作为一个MQTT客户端订阅所有设备数据相关的主题。收到数据后进行解析、校验、业务逻辑处理如超限告警、数据清洗然后持久化到数据库如MySQL/PostgreSQL用于关系型数据InfluxDB或TDengine用于时序数据同时通过WebSocket或Server-Sent Events将实时数据推送给前端。数据展示与交互层前端Web应用。使用Vue.js或React等框架结合ECharts、AntV等图表库构建实时监控仪表盘。它通过WebSocket与SpringBoot服务保持长连接接收实时数据流并动态更新图表。这个架构的优势在于每一层都可以独立扩展和升级。比如设备增加了只需要让新设备按规范发布消息到MQTT存储压力大了可以单独优化数据库或引入缓存前端展示需要更换大屏框架也不影响后端逻辑。2.2 核心组件选型与考量MQTT Broker选型EMQX市面上Broker不少比如Mosquitto轻量、HiveMQ商用、NanoMQ边缘。选择EMQX 5.x版本主要看中它1对MQTT 5.0协议的完整支持带来了更好的错误处理和会话管理2强大的规则引擎可以在Broker端直接对数据进行简单的过滤、转换甚至写入数据库减轻后端压力3易于监控和管理Web控制台能清晰看到连接数、消息吞吐量。对于实验室百台设备以内的规模单机部署完全够用未来如需扩展其集群方案也很成熟。SpringBoot与MQTT客户端库SpringBoot版本选择当前稳定的2.7.x或3.x系列。集成MQTT客户端我们选用org.springframework.integration:spring-integration-mqtt。这个库是Spring Integration项目的一部分它提供了更“Spring”风格的配置方式将MQTT连接抽象为MessageChannel通过注解ServiceActivator就能处理消息与Spring生态无缝集成。相比直接使用Paho客户端它简化了连接管理、重连等繁琐逻辑。数据库选型MySQL InfluxDB数据存储分两类。设备元数据、用户信息、告警配置等用MySQL。而设备产生的时序数据如温度、压力、电压等带时间戳的测点数据则强烈推荐使用时序数据库。这里选InfluxDB因为它专为时序数据优化写入和按时间范围查询的速度极快压缩率高自带类SQL的查询语言Flux也容易上手。对于监控场景下“高写入、少更新、按时间聚合查询”的需求它比关系型数据库合适得多。前端实时通信WebSocket要实现前端图表毫秒级更新HTTP轮询和长轮询都不可取。SpringBoot通过spring-boot-starter-websocket可以轻松集成WebSocket。我们会在后端维护一个全局的WebSocketSession池当SpringBoot的MQTT客户端收到新数据并处理完后立即将数据广播给所有在线的WebSocket会话前端监听消息并更新DOM。3. 核心细节解析与实操要点3.1 MQTT主题设计与消息规范主题设计是MQTT应用的基础好的设计能让系统清晰且易于扩展。我们采用分层结构lab/device/${deviceId}/sensor/${sensorType}例如ID为FURNACE_01的加热炉的温度传感器数据主题为lab/device/FURNACE_01/sensor/temperature。这种结构便于订阅比如订阅lab/device//sensor/可以收到所有设备所有传感器的数据订阅lab/device/FURNACE_01/sensor/可以收到该设备所有传感器的数据。消息体统一采用JSON格式包含必要字段{ deviceId: FURNACE_01, sensorType: temperature, value: 356.8, unit: °C, timestamp: 1715589123456, status: normal // 可选如 normal, warning, fault }注意timestamp建议使用设备采集时的UTC时间戳毫秒避免因网络传输延迟导致服务端时间不准。如果设备无法提供则在数据采集程序中打上时间戳。实操心得在主题中加入版本号是个好习惯例如lab/v1/device/...为未来协议升级留有余地。消息体不要过于庞大单个消息尽量只包含一个测点的数据或者同一设备同一时刻的一组相关数据避免因某个传感器故障导致整包数据丢失。3.2 SpringBoot集成MQTT客户端详解在SpringBoot中配置MQTT客户端我们主要利用MqttPahoClientFactory和MessageProducer。首先在pom.xml引入依赖dependency groupIdorg.springframework.integration/groupId artifactIdspring-integration-mqtt/artifactId /dependency dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.2.5/version /dependency接着通过Java Config方式配置工厂和入站/出站通道适配器Configuration public class MqttConfig { Value(${mqtt.broker-url}) private String brokerUrl; Value(${mqtt.client-id}) private String clientId; Bean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory new DefaultMqttPahoClientFactory(); MqttConnectOptions options new MqttConnectOptions(); options.setServerURIs(new String[]{brokerUrl}); options.setCleanSession(true); // 设为true不持久化会话 options.setAutomaticReconnect(true); // 关键启用自动重连 options.setConnectionTimeout(10); // 如果有用户名密码 // options.setUserName(admin); // options.setPassword(password.toCharArray()); factory.setConnectionOptions(options); return factory; } // 入站通道适配器用于订阅消息 Bean public MessageChannel mqttInputChannel() { return new DirectChannel(); } Bean public MessageProducer inbound() { MqttPahoMessageDrivenChannelAdapter adapter new MqttPahoMessageDrivenChannelAdapter(brokerUrl, clientId _inbound, mqttClientFactory(), lab/device//sensor/); adapter.setCompletionTimeout(5000); adapter.setConverter(new DefaultPahoMessageConverter()); adapter.setQos(1); // 服务质量设为1确保消息至少到达一次 adapter.setOutputChannel(mqttInputChannel()); return adapter; } // 出站通道适配器用于发布消息如下发控制指令 Bean ServiceActivator(inputChannel mqttOutboundChannel) public MessageHandler outbound() { MqttPahoMessageHandler handler new MqttPahoMessageHandler(brokerUrl, clientId _outbound, mqttClientFactory()); handler.setAsync(true); handler.setDefaultTopic(lab/device/control); handler.setDefaultQos(1); return handler; } }提示clientId需要保证唯一性通常用“服务名实例标识”来构造。setAutomaticReconnect(true)对于生产环境至关重要它能应对网络闪断。QoS设置为1是可靠性和性能的平衡对于实验室监控确保数据不丢失比极致低延迟更重要。3.3 时序数据存储与InfluxDB集成对于源源不断的传感器数据我们需要高效存储。在SpringBoot中集成InfluxDB 2.x首先引入客户端依赖dependency groupIdcom.influxdb/groupId artifactIdinfluxdb-client-java/artifactId version6.10.0/version /dependency编写一个服务类负责将处理后的传感器数据写入InfluxDBService public class InfluxDbService { Value(${influxdb.url}) private String url; Value(${influxdb.token}) private String token; Value(${influxdb.org}) private String org; Value(${influxdb.bucket}) private String bucket; private InfluxDBClient influxDBClient; PostConstruct public void init() { this.influxDBClient InfluxDBClientFactory.create(url, token.toCharArray(), org, bucket); } public void writeSensorData(String deviceId, String sensorType, Double value, String unit, Instant timestamp) { try (WriteApi writeApi influxDBClient.getWriteApi()) { Point point Point.measurement(sensor_data) .addTag(device_id, deviceId) .addTag(sensor_type, sensorType) .addField(value, value) .addField(unit, unit) .time(timestamp, WritePrecision.MS); writeApi.writePoint(point); } catch (Exception e) { // 记录日志并考虑将失败数据暂存到本地队列或数据库后续重试 log.error(写入InfluxDB失败: deviceId{}, sensorType{}, deviceId, sensorType, e); } } PreDestroy public void close() { if (influxDBClient ! null) { influxDBClient.close(); } } }InfluxDB的数据模型基于Measurement类似表、Tag索引字段如device_id、Field值字段如value和Time。Tag用于高效过滤和分组Field用于存储实际数值。上述代码将每条数据作为一个Point写入。注意事项写入InfluxDB时时间戳timestamp的一致性非常重要。如果多个设备时间不同步会导致查询时数据错乱。最佳实践是使用一个可靠的时间源如NTP服务器为所有设备和服务对时或者如前所述在数据采集端统一使用UTC时间戳。4. 实操过程与核心环节实现4.1 从设备到Broker数据采集端实现数据采集端因设备而异这里以常见的“边缘网关”模式为例。假设我们有一台运行Linux的工控机通过RS485连接多个温湿度传感器。我们可以用Python编写采集程序使用paho-mqtt库。首先读取串口数据使用pyserial库解析出温湿度值。然后连接EMQX Broker并发布消息import paho.mqtt.client as mqtt import json import time broker 192.168.1.100 port 1883 client_id fgateway-python-{int(time.time())} def on_connect(client, userdata, flags, rc): if rc 0: print(Connected to MQTT Broker!) else: print(fFailed to connect, return code {rc}) client mqtt.Client(client_id) client.on_connect on_connect client.connect(broker, port) # 模拟从串口读取并解析数据 def read_sensor_from_serial(): # 这里省略具体的串口读取和协议解析代码 # 假设解析后得到 device_id, temp, humidity return SENSOR_01, 25.3, 65.2 while True: device_id, temperature, humidity read_sensor_from_serial() current_timestamp int(time.time() * 1000) # 毫秒时间戳 # 发布温度数据 temp_topic flab/device/{device_id}/sensor/temperature temp_payload json.dumps({ deviceId: device_id, sensorType: temperature, value: temperature, unit: °C, timestamp: current_timestamp }) client.publish(temp_topic, temp_payload, qos1) # 发布湿度数据 humi_topic flab/device/{device_id}/sensor/humidity humi_payload json.dumps({ deviceId: device_id, sensorType: humidity, value: humidity, unit: %RH, timestamp: current_timestamp }) client.publish(humi_topic, humi_payload, qos1) time.sleep(5) # 每5秒采集一次关键点采集端需要做好异常处理和重连机制。网络中断或Broker重启时paho-mqtt客户端需要能够自动重连。此外对于重要数据可以考虑在本地进行缓存待网络恢复后重发但这会增加采集端的复杂性可根据数据重要性权衡。4.2 SpringBoot服务端消息接收、处理与广播在SpringBoot中我们通过服务激活器来处理MQTT订阅到的消息。Service public class MqttMessageService { Autowired private InfluxDbService influxDbService; Autowired private SimpMessagingTemplate messagingTemplate; // 用于WebSocket广播 ServiceActivator(inputChannel mqttInputChannel) public void handleMessage(Message? message) { String topic (String) message.getHeaders().get(mqtt_receivedTopic); String payload (String) message.getPayload(); try { // 1. 解析JSON ObjectMapper mapper new ObjectMapper(); SensorData sensorData mapper.readValue(payload, SensorData.class); // 2. 数据校验例如值域范围 if (!isValid(sensorData)) { log.warn(收到无效数据: {}, payload); return; } // 3. 持久化到时序数据库 Instant instant Instant.ofEpochMilli(sensorData.getTimestamp()); influxDbService.writeSensorData( sensorData.getDeviceId(), sensorData.getSensorType(), sensorData.getValue(), sensorData.getUnit(), instant ); // 4. 实时告警判断例如温度超限 checkAndAlert(sensorData); // 5. 通过WebSocket广播实时数据给前端 MapString, Object wsMessage new HashMap(); wsMessage.put(deviceId, sensorData.getDeviceId()); wsMessage.put(sensorType, sensorData.getSensorType()); wsMessage.put(value, sensorData.getValue()); wsMessage.put(timestamp, sensorData.getTimestamp()); // 广播到前端订阅的频道例如 /topic/realtime-data messagingTemplate.convertAndSend(/topic/realtime-data, wsMessage); log.debug(已处理数据: {}, sensorData); } catch (JsonProcessingException e) { log.error(MQTT消息JSON解析失败: {}, payload, e); } catch (Exception e) { log.error(处理MQTT消息时发生未知错误, e); } } private boolean isValid(SensorData data) { // 简单的校验逻辑 if (data.getValue() null || data.getDeviceId() null) { return false; } // 可根据传感器类型添加特定范围校验 return true; } private void checkAndAlert(SensorData data) { // 从数据库或缓存读取该设备的告警阈值配置 // 如果 data.getValue() 超过阈值则触发告警动作 // 例如记录告警日志、发送邮件/短信、向前端推送告警消息 } }这个handleMessage方法是整个后端数据流的枢纽。它必须是线程安全的因为MQTT消息是并发到达的。这里使用了Spring Integration的ServiceActivator注解它会自动从mqttInputChannel拉取消息并调用此方法。实操心得在处理消息的链路上要考虑性能。如果设备数据量非常大每秒上千条上述同步处理、写库、广播的流程可能成为瓶颈。此时可以考虑引入消息队列如Kafka、RocketMQ进行削峰填谷或者使用Async注解将写库和广播操作异步化但要注意事务和顺序性问题。4.3 前端实时监控大屏构建前端使用Vue3配合vue-echarts和SockJS-client、stompjs用于WebSocket。首先建立与后端的WebSocket连接import { Client } from stomp/stompjs; import SockJS from sockjs-client; const wsUrl http://${location.host}/ws-endpoint; const client new Client({ webSocketFactory: () new SockJS(wsUrl), reconnectDelay: 5000, heartbeatIncoming: 4000, heartbeatOutgoing: 4000, }); client.onConnect (frame) { console.log(WebSocket连接成功); // 订阅后端广播实时数据的频道 client.subscribe(/topic/realtime-data, (message) { const data JSON.parse(message.body); // 更新对应的图表数据 updateChart(data); }); }; client.activate();然后使用ECharts绘制实时曲线。关键点在于管理图表的数据队列避免数据点过多导致内存和渲染问题。// 假设每个设备每个传感器对应一个图表实例 const chartDataMap new Map(); // key: ${deviceId}-${sensorType} function updateChart(incomingData) { const key ${incomingData.deviceId}-${incomingData.sensorType}; if (!chartDataMap.has(key)) { initChart(key); // 初始化图表 } const chartInstance chartDataMap.get(key); const option chartInstance.getOption(); // 获取原有的数据序列 const seriesData option.series[0].data; // 添加新数据点 [时间戳, 值] seriesData.push([incomingData.timestamp, incomingData.value]); // 限制数据点数量例如只保留最近1000个点 const MAX_POINTS 1000; if (seriesData.length MAX_POINTS) { seriesData.shift(); } // 更新x轴的时间范围使其滚动 option.xAxis[0].data seriesData.map(item new Date(item[0]).toLocaleTimeString()); // 这里简化处理实际应更新data option.series[0].data seriesData.map(item item[1]); chartInstance.setOption(option); }注意事项前端频繁更新DOM如图表是性能敏感操作。一定要做好防抖或节流特别是当数据速率很快时。另外当打开多个监控页面时每个页面都会建立一个WebSocket连接后端需要能管理这些连接。前端在页面卸载时务必记得调用client.deactivate()关闭连接释放资源。5. 系统部署、调优与安全考量5.1 服务端部署与配置建议使用Docker Compose来编排整个后端服务这样部署和迁移都方便。一个简单的docker-compose.yml示例如下version: 3.8 services: emqx: image: emqx/emqx:5.4 container_name: emqx ports: - 1883:1883 # MQTT TCP端口 - 8083:8083 # MQTT WebSocket端口供前端直连测试用 - 18083:18083 # EMQX Dashboard管理界面 environment: - EMQX_LOG__LEVELwarning volumes: - ./emqx_data:/opt/emqx/data networks: - lab-net influxdb: image: influxdb:2.7 container_name: influxdb ports: - 8086:8086 environment: - DOCKER_INFLUXDB_INIT_MODEsetup - DOCKER_INFLUXDB_INIT_USERNAMEadmin - DOCKER_INFLUXDB_INIT_PASSWORDyour_secure_password - DOCKER_INFLUXDB_INIT_ORGmy-lab - DOCKER_INFLUXDB_INIT_BUCKETsensor_data - DOCKER_INFLUXDB_INIT_ADMIN_TOKENyour_super_secret_token volumes: - ./influxdb_data:/var/lib/influxdb2 networks: - lab-net mysql: image: mysql:8.0 container_name: mysql ports: - 3306:3306 environment: - MYSQL_ROOT_PASSWORDyour_root_password - MYSQL_DATABASElab_monitor volumes: - ./mysql_data:/var/lib/mysql networks: - lab-net backend: build: ./backend # 指向你的SpringBoot项目Dockerfile所在目录 container_name: springboot-backend ports: - 8080:8080 environment: - SPRING_PROFILES_ACTIVEprod - MQTT_BROKER_URLtcp://emqx:1883 - INFLUXDB_URLhttp://influxdb:8086 - INFLUXDB_TOKENyour_super_secret_token depends_on: - emqx - influxdb - mysql networks: - lab-net networks: lab-net: driver: bridge部署时关键是将服务间的连接地址从localhost改为Docker Compose中定义的服务名如emqx,influxdb。SpringBoot的application-prod.yml配置文件需要相应调整。5.2 性能调优与稳定性保障EMQX调优连接数默认配置支持大量连接但需要根据服务器内存调整。关注emqx.conf中的zone.external.max_connections。会话与消息持久化如果消息非常重要可以启用EMQX的持久化功能将消息和会话存储到外部数据库如MySQL、PostgreSQL防止Broker重启导致数据丢失。但这会牺牲一部分性能。规则引擎对于简单的数据过滤或转发可以利用EMQX的规则引擎在Broker端完成减少后端服务的压力。例如将所有温度超过100度的消息单独写入一个主题供告警服务订阅。SpringBoot服务调优连接池确保数据库MySQL、InfluxDB客户端使用了连接池如HikariCP并合理配置最大连接数。异步处理如前所述对于耗时的操作如复杂的告警计算、写入数据库使用Async配合线程池避免阻塞MQTT消息监听线程。JVM参数在生产环境部署时根据服务器内存合理设置JVM堆内存-Xms和-Xmx并启用GC日志以便排查问题。前端优化数据聚合对于历史曲线查询不要一次性拉取所有原始数据。后端应提供聚合查询接口如查询过去24小时按每5分钟求平均值前端只请求聚合后的数据大幅减少传输量和前端渲染压力。虚拟滚动/分页如果监控列表设备很多采用虚拟滚动技术只渲染可视区域内的设备卡片。5.3 安全加固措施实验室系统虽在内网安全也不容忽视。MQTT通信安全禁用匿名访问在EMQX中配置allow_anonymous false强制所有连接必须提供用户名密码或客户端证书。使用ACL配置访问控制列表限制每个客户端只能订阅和发布特定的主题。例如数据采集客户端只能发布到lab/device/${自身ID}/sensor/而不能订阅其他设备的数据。启用TLS/SSL使用mqtts://协议和端口8883对传输层进行加密。需要为EMQX配置证书。修改默认端口将默认的1883端口改为其他不常见的端口减少被扫描攻击的风险。应用层安全输入校验SpringBoot服务端对收到的MQTT消息必须进行严格的JSON解析和字段校验防止注入攻击或畸形数据导致服务崩溃。API鉴权前端与后端SpringBoot的REST API交互如查询历史数据、修改配置需要使用JWT等机制进行鉴权。WebSocket连接验证在建立WebSocket连接时可以要求前端提供Token后端在握手阶段进行验证。数据库安全最小权限原则为SpringBoot应用创建独立的数据库用户只授予必要的读写权限不要使用root账户。InfluxDB令牌管理妥善保管InfluxDB的初始化Token并为不同服务创建具有不同权限的独立API Token。6. 常见问题与排查技巧实录在实际部署和运行中肯定会遇到各种问题。这里记录几个我踩过的坑和解决方法。6.1 MQTT连接与通信问题问题1设备端频繁断开连接日志显示“Connection lost”或“无法连接到Broker”。排查检查网络连通性在设备端ping一下Broker的IP地址。检查防火墙确保Broker所在服务器的1883端口或自定义端口已开放。检查Broker状态登录EMQX Dashboard查看Broker是否正常运行资源CPU、内存是否吃紧。检查客户端ID冲突MQTT协议要求同一Broker上连接的客户端ID必须唯一。如果两个设备用了相同的clientId后连接的会踢掉先连接的。确保采集端代码中的clientId是动态或唯一的。解决在客户端代码中务必设置setAutomaticReconnect(true)和合理的重连间隔。对于clientId可以用“设备型号MAC地址随机数”的组合来生成。问题2SpringBoot服务订阅不到某些主题的消息。排查检查主题匹配确认SpringBoot中MqttPahoMessageDrivenChannelAdapter设置的订阅主题通配符是否正确。例如设备发布到lab/device/01/temp但SpringBoot订阅的是lab/device//temperature当然收不到。检查QoS级别如果设备发布时QoS0而服务端网络不稳定消息可能丢失。确保重要数据使用QoS1或2。查看EMQX Dashboard在“监控 - 主题监控”中查看目标主题是否有消息流入以及订阅者列表里是否有你的SpringBoot客户端。解决在代码中打印出收到的消息主题和内容与设备发布的消息进行比对。使用EMQX的“WebSocket客户端”工具手动发布一条消息测试服务端能否收到可以快速定位是发布端问题还是订阅端问题。6.2 数据存储与查询性能问题问题随着数据量增长查询历史曲线越来越慢。排查检查InfluxDB数据保留策略默认的保留策略可能是无限期导致数据量无限增长。需要根据业务需求设置合理的RPRetention Policy例如只保留30天的原始数据。检查查询语句是否在查询中使用了无法利用索引的过滤条件对于InfluxDB对tag的过滤效率远高于对field的过滤。检查硬件资源InfluxDB的写入和查询都比较消耗IO和内存监控服务器磁盘IO和内存使用情况。解决设置数据保留策略CREATE RETENTION POLICY 30_days ON my-lab DURATION 30d REPLICATION 1 DEFAULT。优化查询前端查询历史数据时一定要带上时间范围限制并且尽量按tag如device_id过滤。对于展示长时间跨度的图表务必使用聚合函数如mean()求平均值下采样而不是查询所有原始数据点。考虑分库分表如果设备数量极多可以考虑按设备类型或区域将数据写入不同的Measurement甚至不同的Bucket。6.3 前端实时数据显示异常问题WebSocket连接不稳定图表数据时有时无或延迟很高。排查检查浏览器控制台查看WebSocket连接是否有错误是否在不断重连。检查网络环境特别是如果前端通过公网访问内网服务网络延迟和稳定性是主要问题。检查后端广播逻辑在SpringBoot服务中确认SimpMessagingTemplate.convertAndSend方法是否被成功调用是否有异常被吞掉。检查STOMP心跳STOMP协议依赖心跳保活。如果网络设备如某些企业防火墙会关闭长时间空闲的TCP连接需要调整心跳间隔。解决在前端WebSocket客户端配置中合理设置reconnectDelay和心跳参数heartbeatIncoming/heartbeatOutgoing。在后端SpringBoot的WebSocket配置中也配置相应的心跳。Configuration EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { Override public void configureWebSocketTransport(WebSocketTransportRegistration registration) { registration.setSendTimeLimit(15 * 1000).setSendBufferSizeLimit(512 * 1024); } Override public void configureClientInboundChannel(ChannelRegistration registration) { registration.taskExecutor().corePoolSize(10).maxPoolSize(20); } // ... 其他配置 }在前端图表更新逻辑中加入“数据陈旧”判断。如果超过一定时间如10秒没有收到新数据则在界面上显示“连接中断”或“数据延迟”的提示。这套系统从搭建到稳定运行是一个不断迭代和优化的过程。最开始可能只关注功能实现跑通流程。随着设备增多、数据量变大就要开始关注架构的弹性、性能和可维护性。比如可以考虑将告警逻辑独立成一个微服务将数据持久化操作放入消息队列异步处理甚至引入流处理框架如Flink进行实时数据分析。
返回列表