某家电企业用Java搭建远程监控系统年省百万运维成本详解
写在前面
去年,我们团队接手了一个相当棘手的活儿——给一家做智能家电的企业做远程监控系统。这家公司每年光运维成本就烧掉几百万,现场工程师跑断腿,用户投诉不断。老板一拍桌子:”必须搞定!”
经过半年多的高强度开发,系统上线后第一个季度,运维成本直接砍了70%。今天就把整个过程掰开了揉碎了讲给你听,从0到1帮你搭一套靠谱的系统。
一、问题出在哪——不先搞清楚痛点就开干,等于白干
1.1 改造前的真实场景
先说几个现场情况,看完你就知道为什么必须做这个系统:
- 故障发现滞后:用户家里空调出问题了,要等用户自己打电话投诉,我们才知道。平均从故障发生到发现问题,要47小时。
- 远程调试不可能:工程师每次上门,要带一堆工具,而且很多时候到了现场才发现参数不对,白跑一趟。
- 数据孤岛:生产数据、售后数据、用户反馈数据各自为政,完全对不上账。
1.2 一个真实案例让我印象深刻
有个用户家的热水器报错了,客服接到电话后第一反应是:”您能描述一下故障现象吗?”用户说”不热”,客服又问”具体多热?”……来回沟通了40分钟,最后派工程师上门,到了才发现是Wi-Fi模块断连了,数据传不上来。
如果系统能在数据断流的第一时间自动告警,工程师提前知道要带Wi-Fi模块备件,整个流程能从40分钟沟通+一天上门,缩短到5分钟告警+1小时精准上门。
二、整体架构设计——顶层设计决定了系统的天花板
2.1 我们选的架构模式
没有用那种花里胡哨的微服务全拆分,而是根据家电行业的特点,选了一个分层架构 + 事件驱动的混合方案:
┌─────────────────────────────────────────────────────────┐
│ 接入层 (Gateway Layer) │
│ MQTT Broker + HTTP API Gateway │
├─────────────────────────────────────────────────────────┤
│ 业务处理层 (Processing Layer) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ 规则引擎 │ │ 数据清洗 │ │ 实时计算 │ │
│ └──────────┘ └──────────┘ └──────────┘ │
├─────────────────────────────────────────────────────────┤
│ 数据存储层 (Storage Layer) │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │ Timescale│ │ Redis │ │ MySQL │ │ ES │ │
│ │ DB │ │ │ │ │ │ │ │
│ └─────────┘ └─────────┘ └─────────┘ └─────────┘ │
├─────────────────────────────────────────────────────────┤
│ 应用服务层 (Service Layer) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ 设备管理 │ │ 告警中心 │ │ 数据分析 │ │
│ └──────────┘ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────┘
2.2 核心设计决策
第一,设备数据流向必须清晰。 传感器 → 网关 → 云平台 → 存储 → 应用,每个环节的数据格式和协议都要明确,不能糊弄。
第二,实时性和可靠性不能两害相权。 家电监控最怕数据丢,所以我们用了消息队列做缓冲,关键告警信息必须100%到达。
第三,可扩展性要提前规划。 一开始设备量少,可能1万台就够了,但系统要能支撑到100万台甚至1000万台。
三、传感器数据采集——这是整个系统的根基
3.1 传感器数据的特性
家电传感器数据有几个典型特点:
- 高频率:温度传感器每秒可能上报几次
- 低价值密度:大部分数据是”正常值”,真正有价值的只是异常点
- 间歇性:设备不是一直在线,用户可能出门、断电、Wi-Fi波动
3.2 数据采集端的设计
我们给每台设备配了一个轻量级的数据采集SDK(Java实现),部署在设备的网关或边缘计算模块上:
/**
* 传感器数据采集SDK核心类
* 负责从本地传感器读取数据,缓存、过滤、上报
*/
@Component
@Slf4j
public class SensorCollector implements InitializingBean {
@Autowired
private MqttClientWrapper mqttClient;
@Autowired
private DataFilterService dataFilterService;
@Autowired
private DeviceConfigService configService;
/**
* 数据缓存队列 - 网络中断时先用本地队列存着
*/
private final BlockingQueue<SensorData> localBuffer =
new ArrayBlockingQueue<>(10000);
/**
* 上报间隔配置(默认30秒)
*/
private volatile int reportIntervalSeconds = 30;
/**
* 传感器数据模板
*/
@Data
@Builder
public static class SensorData {
private String deviceId;
private String sensorType; // 温度/湿度/功耗/振动等
private Double value; // 传感器数值
private Integer quality; // 数据质量 (0-100)
private Long timestamp; // 采集时间戳(毫秒)
private String unit; // 单位
private String firmwareVersion; // 固件版本
}
/**
* 定时采集任务
*/
@Scheduled(fixedRate = 1000)
public void collectSensorData() {
try {
// 1. 从本地传感器读取原始数据
List<SensorData> rawDataList = readRawSensorData();
// 2. 数据过滤 - 去除无效数据
List<SensorData> filteredData = rawDataList.stream()
.filter(this::isDataValid)
.collect(Collectors.toList());
// 3. 数据压缩 - 对连续相同值做去重
List<SensorData> compressedData =
compressSameValueData(filteredData);
// 4. 逐条上报或批量上报
compressedData.forEach(this::uploadToCloud);
} catch (Exception e) {
log.error("传感器采集异常, deviceId={}",
getDeviceId(), e);
}
}
/**
* 读取原始传感器数据
*/
private List<SensorData> readRawSensorData() {
// 实际项目中这里是调用硬件接口读取
// 我们用一个模拟实现来说明逻辑
List<SensorData> data = new ArrayList<>();
// 模拟读取温度传感器
SensorData temperature = SensorData.builder()
.deviceId(getDeviceId())
.sensorType("TEMPERATURE")
.value(getCurrentTemperature())
.quality(95)
.timestamp(System.currentTimeMillis())
.unit("°C")
.firmwareVersion("v2.3.1")
.build();
data.add(temperature);
// 模拟读取湿度传感器
SensorData humidity = SensorData.builder()
.deviceId(getDeviceId())
.sensorType("HUMIDITY")
.value(getCurrentHumidity())
.quality(92)
.timestamp(System.currentTimeMillis())
.unit("%RH")
.firmwareVersion("v2.3.1")
.build();
data.add(humidity);
// 模拟读取功耗传感器
SensorData power = SensorData.builder()
.deviceId(getDeviceId())
.sensorType("POWER")
.value(getCurrentPower())
.quality(88)
.timestamp(System.currentTimeMillis())
.unit("W")
.firmwareVersion("v2.3.1")
.build();
data.add(power);
return data;
}
/**
* 数据质量判断 - 过滤无效数据
*/
private boolean isDataValid(SensorData data) {
if (data.getValue() == null) {
return false;
}
// 温度范围校验
if ("TEMPERATURE".equals(data.getSensorType())) {
return data.getValue() >= -40 && data.getValue() <= 150;
}
// 湿度范围校验
if ("HUMIDITY".equals(data.getSensorType())) {
return data.getValue() >= 0 && data.getValue() <= 100;
}
return true;
}
/**
* 对连续相同值做数据压缩,减少上传量
*/
private List<SensorData> compressSameValueData(
List<SensorData> dataList) {
if (dataList.isEmpty()) {
return dataList;
}
List<SensorData> result = new ArrayList<>();
SensorData lastData = dataList.get(0);
int consecutiveCount = 1;
for (int i = 1; i < dataList.size(); i++) {
SensorData current = dataList.get(i);
// 同类型传感器,数值差异小于阈值(温度0.5度,湿度2%)
boolean isSameValue = isValueSimilar(lastData, current);
if (isSameValue && consecutiveCount < 10) {
// 合并,计数+1
consecutiveCount++;
lastData = current;
} else {
// 值变化了或者连续相同次数太多,添加并更新
if (consecutiveCount > 1) {
// 批量压缩,减少传输次数
lastData.setQuality(
lastData.getQuality() - consecutiveCount);
}
result.add(lastData);
lastData = current;
consecutiveCount = 1;
}
}
result.add(lastData);
return result;
}
private boolean isValueSimilar(SensorData a, SensorData b) {
if (!a.getSensorType().equals(b.getSensorType())) {
return false;
}
double threshold = "TEMPERATURE".equals(a.getSensorType())
? 0.5 : 2.0;
return Math.abs(a.getValue() - b.getValue()) < threshold;
}
/**
* 上传到云端 - 支持本地缓冲
*/
private void uploadToCloud(SensorData data) {
try {
// 尝试直接上传
if (mqttClient.isConnected()) {
String topic = buildUploadTopic(data.getDeviceId());
String payload = JSON.toJSONString(data);
mqttClient.publish(topic, payload.getBytes(), 1);
} else {
// 网络中断,存入本地队列
localBuffer.offer(data);
log.warn("网络中断,数据已缓冲, deviceId={}, sensorType={}",
data.getDeviceId(), data.getSensorType());
}
} catch (Exception e) {
// 上传失败也存入本地队列,保证不丢数据
localBuffer.offer(data);
log.error("数据上传失败,已缓冲, deviceId={}",
data.getDeviceId(), e);
}
}
/**
* 定期flush本地缓冲队列
*/
@Scheduled(fixedDelay = 5000)
public void flushLocalBuffer() {
if (mqttClient.isConnected() && !localBuffer.isEmpty()) {
int flushed = 0;
while (!localBuffer.isEmpty() && flushed < 100) {
SensorData data = localBuffer.poll();
if (data != null) {
try {
String topic = buildUploadTopic(data.getDeviceId());
String payload = JSON.toJSONString(data);
mqttClient.publish(topic, payload.getBytes(), 1);
flushed++;
} catch (Exception e) {
// 重新放回队列
localBuffer.offer(data);
break;
}
}
}
if (flushed > 0) {
log.info("本地缓冲刷新完成, deviceId={}, flushed={}",
getDeviceId(), flushed);
}
}
}
private String buildUploadTopic(String deviceId) {
return String.format("sensor/%s/data", deviceId);
}
@Override
public void afterPropertiesSet() {
// 启动时恢复本地缓冲中的数据
recoverLocalBuffer();
}
}
3.3 数据采集的几个关键设计点
点1:本地缓存机制
传感器采集端必须考虑网络不稳定的情况。我们用了ArrayBlockingQueue做本地缓冲,网络断了数据不丢,等网络恢复后自动补传。这个设计救过我们很多次——有台设备在地下室,信号时断时续,但因为本地缓冲机制,数据完整率保持在99.7%以上。
点2:数据压缩 家电的传感器数据有很强的连续性,温度不会一秒内从20度跳到50度。我们对连续相同值做了压缩,把原本需要上传100条数据压缩到只上传5条。数据量直接减少90%以上,传输成本和存储成本都下来了。
点3:数据质量标记
每条数据都带一个quality字段(0-100),反映数据的可信度。传感器老化、接触不良时quality会下降。这样云端在分析时就能知道哪些数据是可信的,哪些需要怀疑。
四、MQTT与HTTP协议对比——为什么我们最终选了MQTT
4.1 两种协议的直观对比
| 对比维度 | MQTT | HTTP |
|---|---|---|
| 协议类型 | 消息协议 | 请求/响应协议 |
| 连接方式 | 长连接(TCP) | 短连接(默认) |
| 数据推送 | 支持服务端主动推送 | 客户端必须主动拉取 |
| 协议开销 | 很小(2字节头部) | 较大(完整HTTP头) |
| 适合场景 | IoT设备、实时性要求高 | 传统Web、API调用 |
| QoS支持 | 3个等级 | 无 |
| 断网处理 | 原生支持遗嘱和保留消息 | 需要自己实现 |
| 实现复杂度 | 中 | 低 |
4.2 一个真实的实验数据
我们做过一个对比测试,1万台家电设备,每台每分钟上报10条传感器数据:
HTTP方案:
- 每台设备每分钟发起600次HTTP请求
- 总请求量:1万台 × 600 = 600万次/分钟
- 平均每次请求开销:约2KB(Headers + Payload)
- 带宽消耗:12GB/分钟
- 服务器连接数峰值:约60万(每分钟新连接)
- 实现难度:低,标准Spring Boot就能做
MQTT方案:
- 1万台设备保持1万个长连接
- 消息频率:每分钟600万条MQTT消息
- 平均每次消息开销:约200字节(MQTT头 + Payload)
- 带宽消耗:约1.2GB/分钟
- 服务器连接数峰值:1万(始终保持)
- 实现难度:中,需要部署MQTT Broker
结论一目了然:MQTT的带宽消耗只有HTTP的1⁄10,连接管理效率高出几百倍。
4.3 我们的MQTT部署方案
/**
* MQTT客户端封装 - 设备端使用
* 基于 Eclipse Paho MQTT Client
*/
@Component
@Slf4j
public class MqttClientWrapper implements InitializingBean {
@Value("${mqtt.broker.url:tcp://iot.example.com:1883}")
private String brokerUrl;
@Value("${mqtt.client.id:device-unknown}")
private String clientId;
@Value("${mqtt.username:}")
private String username;
@Value("${mqtt.password:}")
private String password;
@Value("${mqtt.keep.alive:60}")
private int keepAlive;
@Value("${mqtt.clean.session:true}")
private boolean cleanSession;
@Value("${mqtt.qos:1}")
private int qos;
@Value("${mqtt.retain:false}")
private boolean retain;
private MqttClient mqttClient;
private MqttConnectOptions connectOptions;
/**
* 遗嘱消息 - 设备异常断开时云端自动感知
*/
private static final String WILL_TOPIC = "device/status/#";
private static final String WILL_MESSAGE = "{\"status\":\"offline\",\"reason\":\"unexpected_disconnect\"}";
@Override
public void afterPropertiesSet() {
initClient();
}
private void initClient() {
try {
// 创建客户端实例
this.mqttClient = new MqttClient(
brokerUrl,
clientId,
new MqttPersistence() // 持久化存储,断网重连后恢复
);
// 配置连接选项
this.connectOptions = new MqttConnectOptions();
connectOptions.setBrokerURL(brokerUrl);
connectOptions.setCleanSession(cleanSession);
connectOptions.setKeepAliveInterval(keepAlive);
connectOptions.setAutomaticReconnect(true); // 自动重连
connectOptions.setWill(WILL_TOPIC,
WILL_MESSAGE.getBytes(), qos, retain);
if (StringUtils.isNotBlank(username)) {
connectOptions.setUserName(username);
connectOptions.setPassword(password.toCharArray());
}
// 设置回调
mqttClient.setCallback(new MqttCallbackExtended() {
@Override
public void connectComplete(boolean reconnect, String serverURI) {
log.info("MQTT连接成功, reconnect={}, server={}",
reconnect, serverURI);
// 重新订阅所有主题
subscribeAllTopics();
}
@Override
public void connectionLost(Throwable cause) {
log.error("MQTT连接丢失, 将自动重连", cause);
}
@Override
public void messageArrived(String topic, MqttMessage message) {
handleMessage(topic, message);
}
@Override
public void deliveryComplete(IMqttDeliveryToken token) {
// 消息投递完成回调
}
});
// 建立连接
mqttClient.connect(connectOptions);
log.info("MQTT客户端初始化完成");
} catch (Exception e) {
log.error("MQTT客户端初始化失败", e);
throw new RuntimeException(e);
}
}
/**
* 发布消息到指定主题
*
* @param topic 主题
* @param payload 消息体
* @param qos 服务质量等级
*/
public void publish(String topic, byte[] payload, int qos)
throws MqttException {
if (!isConnected()) {
throw new MqttException(
MqttException.REASON_CODE_CLIENT_NOT_CONNECTED);
}
MqttMessage message = new MqttMessage(payload);
message.setQos(qos);
message.setRetained(false);
mqttClient.publish(topic, message);
log.debug("消息发布成功, topic={}, qos={}", topic, qos);
}
/**
* 订阅主题
*/
public void subscribe(String topic, int qos) throws MqttException {
if (!isConnected()) {
throw new MqttException(
MqttException.REASON_CODE_CLIENT_NOT_CONNECTED);
}
mqttClient.subscribe(topic, qos);
log.info("订阅主题成功, topic={}", topic);
}
/**
* 订阅所有设备主题
*/
private void subscribeAllTopics() {
String[] topics = {
"device/command/#", // 设备指令
"device/config/#", // 配置更新
"device/firmware/#" // 固件升级
};
try {
for (String topic : topics) {
subscribe(topic, 1);
}
} catch (Exception e) {
log.error("订阅主题失败", e);
}
}
public boolean isConnected() {
return mqttClient != null && mqttClient.isConnected();
}
/**
* 处理接收到的消息
*/
private void handleMessage(String topic, MqttMessage message) {
try {
String payload = new String(message.getPayload());
log.debug("收到消息, topic={}, payload={}", topic, payload);
// 根据主题分发到不同处理器
if (topic.startsWith("device/command/")) {
handleDeviceCommand(payload);
} else if (topic.startsWith("device/config/")) {
handleConfigUpdate(payload);
}
} catch (Exception e) {
log.error("处理MQTT消息异常, topic={}", topic, e);
}
}
/**
* 处理设备指令(如远程重启、参数调整等)
*/
private void handleDeviceCommand(String payload) {
// 实际项目中这里调用命令处理器
log.info("处理设备指令: {}", payload);
}
/**
* 处理配置更新
*/
private void handleConfigUpdate(String payload) {
log.info("处理配置更新: {}", payload);
}
/**
* 断开连接
*/
public void disconnect() {
try {
if (mqttClient != null && mqttClient.isConnected()) {
mqttClient.disconnect();
}
} catch (Exception e) {
log.error("断开MQTT连接异常", e);
}
}
}
4.4 MQTT Broker的选择
我们团队内部其实 debate 了很久:
- EMQ X:开源免费,性能强,分布式支持好,但部署复杂度稍高
- HiveMQ:商业版功能强大,社区版有限制
- Mosquitto:轻量,适合小规模,大规模时性能不够
- 阿里云IoT:SaaS方案,省去运维,但有数据合规顾虑
最终我们选了EMQ X Cluster,部署3个节点做高可用,原因很简单:设备量会持续增长,现在1万台,一年后可能10万台,EMQ X的集群方案能平滑扩展。
五、云端平台架构——Java后端的核心实现
5.1 整体技术栈
┌──────────────────────────────────────────────┐
│ Nginx 负载均衡 │
├──────────────────────────────────────────────┤
│ Spring Cloud Gateway (API网关) │
├──────────────────────────────────────────────┤
│ 设备服务 │ 告警服务 │ 数据分析服务 │ 用户服务 │
├──────────────────────────────────────────────┤
│ RocketMQ 消息队列 │
├──────────────────────────────────────────────┤
│ TimescaleDB │ Redis │ MySQL │ ES │
└──────────────────────────────────────────────┘
5.2 设备接入服务——核心中的核心
/**
* 设备接入服务 - 处理MQTT消息和设备管理
* 这是整个系统的核心服务,承载着设备连接和消息接收
*/
@Service
@Slf4j
@RequiredArgsConstructor
public class DeviceAccessService {
private final MqttMessageListener mqttMessageListener;
private final DeviceRepository deviceRepository;
private final DeviceStateService deviceStateService;
private final AlertService alertService;
private final RocketMQTemplate rocketMQTemplate;
/**
* MQTT消息监听 - 接收所有设备的传感器数据
*/
@MqttListener(topic = "sensor/+/data", qos = 1)
public void onSensorDataReceived(String topic, byte[] payload) {
long startTime = System.currentTimeMillis();
try {
// 1. 解析设备ID(从topic中提取)
String deviceId = extractDeviceIdFromTopic(topic);
// 2. 解析传感器数据
SensorDataDTO sensorData = parseSensorData(payload);
sensorData.setDeviceId(deviceId);
sensorData.setReceivedTime(System.currentTimeMillis());
// 3. 验证设备是否合法
validateDevice(deviceId);
// 4. 更新设备状态
deviceStateService.updateDeviceStatus(deviceId, "online");
// 5. 异常检测(实时规则引擎)
checkAnomaly(deviceId, sensorData);
// 6. 写入消息队列(异步处理,提高吞吐)
rocketMQTemplate.asyncSend(
"sensor-data-topic",
sensorData,
new SendCallback() {
@Override
public void onSuccess(SendResult result) {
long cost = System.currentTimeMillis() - startTime;
if (cost > 100) {
log.warn("消息发送耗时过长, cost={}ms, deviceId={}",
cost, deviceId);
}
}
@Override
public void onException(Throwable e) {
log.error("消息发送失败, deviceId={}", deviceId, e);
}
}
);
} catch (Exception e) {
log.error("处理传感器数据异常, topic={}, deviceId={}",
topic, extractDeviceIdFromTopic(topic), e);
}
}
/**
* 从MQTT topic中提取设备ID
* 格式: sensor/{deviceId}/data
*/
private String extractDeviceIdFromTopic(String topic) {
String[] parts = topic.split("/");
if (parts.length >= 2) {
return parts[1];
}
throw new IllegalArgumentException("无效的topic格式: " + topic);
}
/**
* 解析传感器数据JSON
*/
private SensorDataDTO parseSensorData(byte[] payload) {
try {
return JSON.parseObject(new String(payload), SensorDataDTO.class);
} catch (Exception e) {
log.error("传感器数据解析失败, payload={}",
new String(payload), e);
throw new IllegalArgumentException("传感器数据格式错误", e);
}
}
/**
* 设备合法性验证
*/
private void validateDevice(String deviceId) {
DeviceDevice device = deviceRepository.findById(deviceId).orElse(null);
if (device == null) {
log.warn("未知设备尝试接入, deviceId={}", deviceId);
throw new IllegalArgumentException("设备未注册: " + deviceId);
}
if (!"active".equals(device.getStatus())) {
log.warn("设备状态异常, deviceId={}, status={}",
deviceId, device.getStatus());
throw new IllegalStateException("设备状态异常: " + device.getStatus());
}
}
/**
* 实时异常检测 - 基于规则引擎
*/
private void checkAnomaly(String deviceId, SensorDataDTO data) {
// 1. 阈值告警
checkThreshold(deviceId, data);
// 2. 趋势异常检测
checkTrend(deviceId, data);
// 3. 设备离线检测
checkOffline(deviceId);
}
/**
* 阈值告警
*/
private void checkThreshold(String deviceId, SensorDataDTO data) {
// 获取该设备该类型传感器的告警阈值配置
AlarmThreshold threshold = alarmThresholdRepository
.findByDeviceIdAndSensorType(deviceId, data.getSensorType());
if (threshold == null) {
return; // 没有配置阈值,跳过
}
boolean isAlarm = false;
String alarmLevel = null;
if (data.getValue() < threshold.getMinValue()) {
isAlarm = true;
alarmLevel = "WARN";
} else if (data.getValue() > threshold.getMaxValue()) {
isAlarm = true;
alarmLevel = "ERROR";
}
if (isAlarm) {
alertService.createAlert(AlertDTO.builder()
.deviceId(deviceId)
.sensorType(data.getSensorType())
.alarmLevel(alarmLevel)
.message(String.format("传感器数值异常: %s = %s",
data.getSensorType(), data.getValue()))
.value(data.getValue())
.thresholdValue(threshold.getMaxValue())
.build());
log.warn("设备告警, deviceId={}, sensor={}, value={}, threshold={}",
deviceId, data.getSensorType(), data.getValue(),
threshold.getMaxValue());
}
}
/**
* 趋势异常检测 - 检测数值突变
*/
private void checkTrend(String deviceId, SensorDataDTO data) {
// 获取最近5条历史数据
List<SensorDataDTO> history = sensorDataRepository
.findLastNByDeviceIdAndSensorType(deviceId,
data.getSensorType(), 5)
.stream()
.map(this::convertToDTO)
.collect(Collectors.toList());
if (history.size() < 3) {
return; // 数据点不够,不检测
}
// 计算变化率
double recentChange = Math.abs(
data.getValue() - history.get(history.size() - 1).getValue());
// 计算历史平均变化率
double avgChange = history.stream()
.limit(4)
.mapToInt(i -> Math.abs(
history.get(i + 1).getValue() - history.get(i).getValue()))
.average()
.orElse(0.0);
// 如果当前变化率超过平均变化率的3倍,触发告警
if (avgChange > 0 && recentChange > avgChange * 3) {
alertService.createAlert(AlertDTO.builder()
.deviceId(deviceId)
.sensorType(data.getSensorType())
.alarmLevel("WARN")
.message(String.format("数值突变: 当前变化率 %s 超过平均变化率 %s 的3倍",
recentChange, avgChange))
.value(data.getValue())
.build());
log.warn("数值突变告警, deviceId={}, sensor={}, change={}",
deviceId, data.getSensorType(), recentChange);
}
}
/**
* 设备离线检测
*/
private void checkOffline(String deviceId) {
// 检查设备是否在配置的时间窗口内没有上报数据
Long lastReportTime = deviceStateService.getLastReportTime(deviceId);
if (lastReportTime != null) {
long offlineDuration = System.currentTimeMillis() - lastReportTime;
int offlineThreshold = 5 * 60 * 1000; // 5分钟
if (offlineDuration > offlineThreshold) {
deviceStateService.updateDeviceStatus(deviceId, "offline");
alertService.createAlert(AlertDTO.builder()
.deviceId(deviceId)
.alarmLevel("ERROR")
.message(String.format("设备离线超过%d分钟",
offlineDuration / 60000))
.build());
log.warn("设备离线告警, deviceId={}, offlineMinutes={}",
deviceId, offlineDuration / 60000);
}
}
}
// ... 其他方法
}
5.3 传感器数据消费服务——消息队列处理
/**
* 传感器数据消费服务
* 从RocketMQ消费传感器数据,写入TimescaleDB
*/
@Service
@Slf4j
@RocketMQMessageListener(
topic = "sensor-data-topic",
consumerGroup = "sensor-data-consumer-group",
maxReconsumeTimes = 3,
consumeThreadMax = 20,
consumeThreadMin = 5
)
@Component
public class SensorDataConsumer implements RocketMQListener<SensorDataDTO> {
@Autowired
private SensorDataBatchWriter sensorDataBatchWriter;
@Autowired
private DeviceFeatureService deviceFeatureService;
@Autowired
private AlertService alertService;
@Override
public void onMessage(SensorDataDTO data) {
try {
// 1. 批量写入时序数据库
sensorDataBatchWriter.save(data);
// 2. 更新设备特征数据(用于后续分析)
deviceFeatureService.updateDeviceFeature(data);
// 3. 更新设备最后上报时间
deviceStateService.updateLastReportTime(
data.getDeviceId(), data.getReceivedTime());
} catch (Exception e) {
log.error("处理传感器数据异常, deviceId={}",
data.getDeviceId(), e);
// 异常由RocketMQ自动重试,最多3次
}
}
}
六、数据库选型——为什么时序数据库是必须的
6.1 数据库选型对比
| 数据库 | 类型 | 适合场景 | 我们的选择 |
|---|---|---|---|
| TimescaleDB | 时序数据库 | 传感器时序数据 | ✅ 主存储 |
| MySQL | 关系型数据库 | 设备元数据、用户数据 | ✅ 辅助存储 |
| Redis | 内存数据库 | 设备状态缓存、告警实时数据 | ✅ 缓存层 |
| Elasticsearch | 搜索引擎 | 日志检索、告警历史查询 | ✅ 辅助存储 |
6.2 TimescaleDB建表脚本
-- 创建扩展
CREATE EXTENSION IF NOT EXISTS timescaledb;
-- 创建传感器数据超表
CREATE TABLE sensor_data (
time TIMESTAMPTZ NOT NULL,
device_id TEXT NOT NULL,
sensor_type TEXT NOT NULL,
value DOUBLE PRECISION NOT NULL,
quality INTEGER NOT NULL DEFAULT 100,
unit TEXT NOT NULL,
metadata JSONB NULL
);
-- 创建超表(自动分区)
SELECT create_hypertable('sensor_data', 'time');
-- 创建索引
CREATE INDEX ON sensor_data (device_id, sensor_type, time DESC);
CREATE INDEX ON sensor_data (time DESC);
-- 创建策略:自动压缩30天前的数据(节省存储空间)
SELECT add_compression_policy('sensor_data', INTERVAL '30 days');
-- 创建策略:自动删除180天前的数据(合规需要)
SELECT add_retention_policy('sensor_data', INTERVAL '180 days');
6.3 为什么选TimescaleDB而不是InfluxDB?
这是一个很实际的选择,我们当时内部也讨论了很久:
TimescaleDB的优势:
- 基于PostgreSQL,SQL查询能力极强,我们团队熟悉度高
- 与现有PostgreSQL生态无缝集成
- 数据压缩效果好,存储空间节省约70%
- 支持复杂查询和 JOIN
InfluxDB的优势:
- 专为时序数据设计,写入性能理论上更高
- 内置查询语言InfluxQL,简单直观
最终选择理由: 我们的场景不仅需要高性能写入,还需要做复杂的关联查询(比如关联设备信息、用户信息、告警记录),TimescaleDB的SQL能力在这方面碾压InfluxDB。而且我们已经有PostgreSQL集群,扩展一个TimescaleDB实例成本几乎为零。
七、并发处理——千万级设备连接的性能挑战
7.1 并发模型选择
MQTT Broker(EMQ X)本身就已经处理了设备连接层面的并发。我们的Java后端主要面对的是消息处理层面的并发。
7.2 我们的并发处理方案
/**
* 设备数据并发处理配置
* 使用线程池 + 消息队列的异步处理架构
*/
@Configuration
@Slf4j
public class DeviceDataProcessingConfig {
/**
* 传感器数据处理线程池
* 核心参数根据设备数量和CPU核数调整
*/
@Bean("sensorDataThreadPool")
public Executor sensorDataThreadPool() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(16); // 核心线程数
executor.setMaxPoolSize(32); // 最大线程数
executor.setQueueCapacity(10000); // 队列容量
executor.setThreadNamePrefix("sensor-data-");
executor.setRejectedExecutionHandler(
new ThreadPoolExecutor.CallerRunsPolicy()); // 拒绝策略:调用者运行
executor.setAwaitTerminationSeconds(60);
executor.setWaitForTasksToCompleteOnShutdown(true);
return executor;
}
/**
* 告警处理线程池
*/
@Bean("alertThreadPool")
public Executor alertThreadPool() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(8);
executor.setMaxPoolSize(16);
executor.setQueueCapacity(5000);
executor.setThreadNamePrefix("alert-");
executor.setRejectedExecutionHandler(
new ThreadPoolExecutor.DiscardPolicy()); // 告警线程池满时丢弃
return executor;
}
/**
* 设备状态更新线程池
*/
@Bean("deviceStateThreadPool")
public Executor deviceStateThreadPool() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(4);
executor.setMaxPoolSize(8);
executor.setQueueCapacity(2000);
executor.setThreadNamePrefix("device-state-");
return executor;
}
}
7.3 关键性能优化点
优化1:批量写入
/**
* 传感器数据批量写入 - 减少数据库连接开销
*/
@Service
@Slf4j
public class SensorDataBatchWriter {
@Autowired
private JdbcTemplate jdbcTemplate;
/**
* 批量保存传感器数据
* 使用INSERT ... VALUES (...), (...), (...) 语法
*/
@Transactional
public void save(List<SensorDataDTO> dataList) {
if (dataList == null || dataList.isEmpty()) {
return;
}
String sql = "INSERT INTO sensor_data " +
"(time, device_id, sensor_type, value, quality, unit) " +
"VALUES (?, ?, ?, ?, ?, ?)";
// 分批写入,每批500条
int batchSize = 500;
for (int i = 0; i < dataList.size(); i += batchSize) {
List<SensorDataDTO> batch = dataList.subList(
i, Math.min(i + batchSize, dataList.size()));
jdbcTemplate.batchUpdate(sql,
new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int j)
throws SQLException {
SensorDataDTO data = batch.get(j);
ps.setTimestamp(1,
new Timestamp(data.getTimestamp()));
ps.setString(2, data.getDeviceId());
ps.setString(3, data.getSensorType());
ps.setDouble(4, data.getValue());
ps.setInt(5, data.getQuality());
ps.setString(6, data.getUnit());
}
@Override
public int getBatchSize() {
return batch.size();
}
});
}
log.debug("批量写入完成, count={}", dataList.size());
}
}
优化2:Redis缓存设备状态
/**
* 设备状态缓存服务
* 使用Redis缓存设备在线状态和最后上报时间
*/
@Service
@Slf4j
public class DeviceStateService {
@Autowired
private RedisTemplate<String, String> redisTemplate;
private static final String STATE_KEY_PREFIX = "device:state:";
private static final String LAST_REPORT_KEY_PREFIX = "device:last_report:";
private static final long STATE_TTL = 10 * 60; // 10分钟过期
/**
* 更新设备在线状态
*/
public void updateDeviceStatus(String deviceId, String status) {
String key = STATE_KEY_PREFIX + deviceId;
redisTemplate.opsForValue().set(key, status,
STATE_TTL, TimeUnit.SECONDS);
// 同时更新最后上报时间
String lastReportKey = LAST_REPORT_KEY_PREFIX + deviceId;
redisTemplate.opsForValue().set(lastReportKey,
String.valueOf(System.currentTimeMillis()),
STATE_TTL, TimeUnit.SECONDS);
}
/**
* 获取设备最后上报时间
*/
public Long getLastReportTime(String deviceId) {
String key = LAST_REPORT_KEY_PREFIX + deviceId;
String value = redisTemplate.opsForValue().get(key);
return value != null ? Long.parseLong(value) : null;
}
/**
* 批量检查离线设备
*/
public List<String> findOfflineDevices() {
long now = System.currentTimeMillis();
long threshold = now - 5 * 60 * 1000; // 5分钟
// 使用Redis SCAN查找所有设备状态key
Set<String> keys = redisTemplate.keys(STATE_KEY_PREFIX + "*");
List<String> offlineDevices = new ArrayList<>();
if (keys != null) {
for (String key : keys) {
String deviceId = key.replace(STATE_KEY_PREFIX, "");
Long lastReport = getLastReportTime(deviceId);
if (lastReport != null && lastReport < threshold) {
offlineDevices.add(deviceId);
}
}
}
return offlineDevices;
}
}
八、常见坑点与排查方法——用血泪换来的经验
8.1 坑点1:设备重复连接导致状态混乱
现象:同一台设备出现多条在线记录,状态时好时坏。
原因:设备重启后重新连接MQTT,旧连接没有完全关闭,导致Broker认为有新设备接入。
排查方法:
-- 查看设备的连接记录
SELECT device_id, status, last_report_time
FROM device_state
WHERE device_id = 'XXX'
ORDER BY last_report_time DESC
LIMIT 10;
-- 查看MQTT会话
SELECT clientid, ipaddress, connected_at
FROM mqtt_client_session
WHERE clientid LIKE '%XXX%';
解决方案:
/**
* MQTT连接去重处理
*/
@Component
@Slf4j
public class MqttConnectionHandler implements MqttConnectHandler {
@Autowired
private DeviceStateService deviceStateService;
@Override
public void handleConnect(MqttClient client, MqttConnectOptions options) {
String clientId = options.getClientId();
String deviceId = extractDeviceId(clientId);
// 检查是否有旧连接
String oldConnection = deviceStateService.getActiveConnection(deviceId);
if (oldConnection != null && !oldConnection.equals(clientId)) {
// 强制断开旧连接
log.warn("发现重复连接,断开旧连接, deviceId={}, oldClient={}, newClient={}",
deviceId, oldConnection, clientId);
forceDisconnect(oldConnection);
}
// 记录新连接
deviceStateService.recordConnection(clientId, deviceId);
}
}
8.2 坑点2:消息堆积导致数据延迟
现象:设备上报数据后,后台延迟数小时才能看到。
排查步骤:
- 检查RocketMQ消费者 lag:
mqadmin consumerProgress - 检查数据库写入速度:监控TimescaleDB写入QPS
- 检查线程池是否满:查看线程池监控指标
解决方案:
# application.yml 调整消费配置
rocketmq:
consumer:
sensor-data:
consumeThreadMin: 20
consumeThreadMax: 64
consumeMaxOffset: 2000 # 最大消费偏移量
8.3 坑点3:MQTT遗嘱消息无法触发
现象:设备异常断电后,云端没有收到离线告警。
原因:遗嘱消息只在TCP连接断开时触发,如果设备是直接断电(不是正常关闭),Broker需要等Keep Alive超时才能检测到。
解决方案:
/**
* 双重离线检测机制
* 1. MQTT遗嘱消息(快速,秒级)
* 2. 超时检测(兜底,分钟级)
*/
@Component
public class DeviceOfflineDetector {
@Autowired
private DeviceStateService deviceStateService;
@Scheduled(fixedRate = 60000) // 每分钟检查一次
public void checkOfflineDevices() {
long threshold = System.currentTimeMillis() - 5 * 60 * 1000;
// 从数据库查询超过5分钟没有上报的设备
List<String> candidates = deviceRepository
.findDevicesWithoutReportBefore(threshold);
for (String deviceId : candidates) {
// 双重确认:检查MQTT连接状态
boolean mqttConnected = mqttClientManager
.isDeviceConnected(deviceId);
if (!mqttConnected) {
// 确认离线,发送告警
deviceStateService.updateDeviceStatus(deviceId, "offline");
alertService.createAlert(AlertDTO.builder()
.deviceId(deviceId)
.alarmLevel("ERROR")
.message("设备离线")
.build());
log.warn("设备确认离线, deviceId={}", deviceId);
}
}
}
}
8.4 坑点4:时序数据写入性能瓶颈
现象:设备量达到5000台以上时,写入延迟明显增加。
排查:
-- 查看TimescaleDB分区情况
SELECT hypertable_name,
count(*) as partition_count,
pg_size_pretty(sum(pg_total_relation_size(
quote_ident(schemaname)||'.'||quote_ident(tablename))
)) as total_size
FROM timescaledb_information.hypertables
GROUP BY hypertable_name;
-- 查看压缩情况
SELECT * FROM timescaledb_information.compression_stats;
解决方案:
- 调整chunk大小(默认1小时,对于高写入场景可以调整为10分钟或30分钟)
- 确保有合适的索引
- 开启自动压缩策略
8.5 坑点5:告警风暴
现象:一次网络波动导致几万条告警同时产生,后台完全瘫痪。
解决方案:告警聚合和去重
/**
* 告警去重服务 - 防止告警风暴
*/
@Service
@Slf4j
public class AlertDeduplicationService {
@Autowired
private RedisTemplate<String, String> redisTemplate;
private static final long DEDUP_WINDOW = 5 * 60; // 5分钟去重窗口
/**
* 判断是否需要发送告警
* 相同设备、相同类型、相同级别的告警在去重窗口内只发一次
*/
public boolean shouldSendAlert(String deviceId, String sensorType,
String alarmLevel) {
String dedupKey = buildDedupKey(deviceId, sensorType, alarmLevel);
Boolean isNew = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", DEDUP_WINDOW, TimeUnit.SECONDS);
return Boolean.TRUE.equals(isNew);
}
private String buildDedupKey(String deviceId, String sensorType,
String alarmLevel) {
return String.format("alert:dedup:%s:%s:%s",
deviceId, sensorType, alarmLevel);
}
}
九、监控与运维——让系统自己”说话”
9.1 核心监控指标
/**
* 系统监控指标注册
*/
@Configuration
public class MonitoringConfig {
@Bean
public MeterRegistryCustomizer<MeterRegistry> metricsCommonTags() {
return registry -> registry.config().commonTags(
"app", "device-monitor",
"env", "production"
);
}
@Bean
public ScheduledExecutorService monitoringScheduler() {
return Executors.newScheduledThreadPool(2);
}
}
/**
* 设备监控指标收集
*/
@Component
@Slf4j
public class DeviceMetricsCollector {
private final MeterRegistry meterRegistry;
public DeviceMetricsCollector(MeterRegistry meterRegistry) {
this.meterRegistry = meterRegistry;
}
/**
* 定时收集设备指标
*/
@Scheduled(fixedRate = 30000) // 每30秒
public void collectDeviceMetrics() {
// 设备在线数量
long onlineCount = deviceStateService.getOnlineDeviceCount();
Gauge.builder("device.online.count", onlineCount)
.register(meterRegistry);
// 每分钟消息接收量
long msgPerMin = mqMetrics.getMessagesPerMinute();
Gauge.builder("mqtt.messages.per_minute", msgPerMin)
.register(meterRegistry);
// 告警数量
long alertCount = alertService.getTodayAlertCount();
Gauge.builder("alert.count.today", alertCount)
.register(meterRegistry);
// 数据处理延迟
long processingLatency = metrics.getAverageProcessingLatency();
Gauge.builder("data.processing.latency_ms", processingLatency)
.register(meterRegistry);
}
}
9.2 日志规范
/**
* 日志规范 - 让问题排查不再是噩梦
*/
@Slf4j
public class DeviceAccessLogger {
/**
* 设备数据接收日志 - 包含完整追踪信息
*/
public void logSensorDataReceived(SensorDataDTO data) {
log.info("[SENSOR_DATA] deviceId={}, sensorType={}, value={},
quality={}, timestamp={}, traceId={}",
data.getDeviceId(),
data.getSensorType(),
data.getValue(),
data.getQuality(),
data.getTimestamp(),
MDC.get("traceId"));
}
/**
* 设备告警日志 - 包含上下文信息
*/
public void logDeviceAlert(String deviceId, AlertDTO alert) {
log.warn("[ALERT] deviceId={}, alarmLevel={}, message={},
value={}, threshold={}, traceId={}",
deviceId,
alert.getAlarmLevel(),
alert.getMessage(),
alert.getValue(),
alert.getThresholdValue(),
MDC.get("traceId"));
}
}
十、成本收益分析——百万运维成本是怎么省下来的
10.1 改造前的成本结构
| 成本项 | 年成本 | 说明 |
|---|---|---|
| 人工巡检 | 120万 | 20个工程师,每月多地出差 |
| 故障处理 | 80万 | 平均每次故障处理2小时+差旅 |
| 用户投诉处理 | 50万 | 客服团队+补偿 |
| 备件浪费 | 30万 | 不知道什么问题,备件带多了 |
| 合计 | 280万 |
10.2 改造后的变化
| 优化项 | 效果 | 年节省 |
|---|---|---|
| 远程诊断替代现场排查 | 80%故障远程解决 | 100万 |
| 提前发现故障 | 故障率降低60% | 50万 |
| 精准备件配送 | 备件浪费减少70% | 20万 |
| 客服压力降低 | 投诉量减少80% | 40万 |
| 合计节省 | 210万 |
10.3 ROI计算
- 系统建设成本:约80万(硬件+开发+部署)
- 年运维成本:约70万(云资源+运维人力)
- 第一年净收益:210万 - 80万 - 70万 = 60万
- 第二年净收益:210万 - 70万 = 140万
- 三年累计收益:约340万
十一、给想搭建类似系统的朋友几点建议
11.1 架构层面
- 不要过早优化:第一阶段先把核心链路跑通,设备量到10万+再考虑分布式。
- 数据分层存储:热数据放Redis/内存,温数据放MySQL,冷数据放TimescaleDB压缩存储。
- 消息队列是必须的:没有MQ,你的系统很难扛住高并发写入。
11.2 代码层面
- 异常处理要彻底:传感器数据格式可能不规范,网络可能随时断,代码要有兜底。
- 日志要规范:每个设备的数据都要有traceId,出问题能追踪到具体设备。
- 监控要提前部署:不要等出问题才发现没有监控。
11.3 运维层面
- 告警分级:不是所有异常都要打电话,区分P0/P1/P2。
- 定期演练:模拟设备大规模离线,检验系统韧性。
- 文档沉淀:把踩过的坑写下来,团队共享。
十二、总结
这套系统从设计到上线,我们团队花了半年时间,踩了不少坑。但回过头看,最值得骄傲的不是技术有多先进,而是实实在在帮企业省下了真金白银。
技术选型上没有追求最热门的,而是选了最合适的:MQTT代替HTTP、TimescaleDB代替MySQL存时序数据、RocketMQ做异步解耦。每个选择都有明确的业务理由,不是拍脑袋决定的。
如果你也在做类似的系统,我的建议是:先跑通最小可行版本,再逐步迭代优化。别一上来就搞微服务全拆分,先把数据链路打通,让设备能上报、数据能入库、异常能告警,这一步就值回票价了。
有问题随时交流,技术这东西,聊着聊着就通透了。
