refactor(iot): 优化事件发布机制并修复状态值解析
Some checks failed
Java CI with Maven / build (11) (push) Has been cancelled
Java CI with Maven / build (17) (push) Has been cancelled
Java CI with Maven / build (8) (push) Has been cancelled

1. IntegrationEventPublisher 只保留设备状态变更事件发布
   - 注释掉 publishPropertyChanged 和 publishEventOccurred 接口
   - RocketMQIntegrationEventPublisher 对应实现改为注释

2. IotDevicePropertyServiceImpl 属性消息发布暂停
   - 注释掉 saveDeviceProperty 中的 publishPropertyMessage 调用
   - 注释掉 publishToIntegrationEventBus 中的实际发布逻辑

3. IotDeviceMessageServiceImpl 新增状态值解析兼容
   - 新增 parseStateValue 方法支持整数和字符串格式状态值
   - 支持 "online"/"offline" 字符串解析

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
lzh
2026-01-21 22:56:13 +08:00
parent 842b40596d
commit fa619710ef
4 changed files with 52 additions and 43 deletions

View File

@@ -14,6 +14,7 @@ import com.viewsh.module.iot.controller.admin.device.vo.message.IotDeviceMessage
import com.viewsh.module.iot.controller.admin.statistics.vo.IotStatisticsDeviceMessageReqVO;
import com.viewsh.module.iot.controller.admin.statistics.vo.IotStatisticsDeviceMessageSummaryByDateRespVO;
import com.viewsh.module.iot.core.enums.IotDeviceMessageMethodEnum;
import com.viewsh.module.iot.core.enums.IotDeviceStateEnum;
import com.viewsh.module.iot.core.mq.message.IotDeviceMessage;
import com.viewsh.module.iot.core.mq.producer.IotDeviceMessageProducer;
import com.viewsh.module.iot.core.util.IotDeviceMessageUtils;
@@ -197,7 +198,9 @@ public class IotDeviceMessageServiceImpl implements IotDeviceMessageService {
String stateStr = IotDeviceMessageUtils.getIdentifier(message);
assert stateStr != null;
Assert.notEmpty(stateStr, "设备状态不能为空");
deviceService.updateDeviceState(device, Integer.valueOf(stateStr));
// 兼容整数和字符串格式的状态值
Integer state = parseStateValue(stateStr);
deviceService.updateDeviceState(device, state);
// TODO 芋艿:子设备的关联
return null;
}
@@ -269,6 +272,34 @@ public class IotDeviceMessageServiceImpl implements IotDeviceMessageService {
});
}
/**
* 解析状态值,支持整数和字符串格式
* <p>
* 支持格式:
* - 整数0=未激活1=在线2=离线
* - 字符串0/inactive=未激活1/online=在线2/offline=离线
*
* @param stateStr 状态字符串
* @return 状态枚举值
*/
private Integer parseStateValue(String stateStr) {
if (stateStr == null) {
return IotDeviceStateEnum.INACTIVE.getState();
}
try {
// 尝试直接解析为整数
return Integer.parseInt(stateStr);
} catch (NumberFormatException e) {
// 字符格式匹配
String lower = stateStr.toLowerCase().trim();
return switch (lower) {
case "1", "online" -> IotDeviceStateEnum.ONLINE.getState();
case "2", "offline" -> IotDeviceStateEnum.OFFLINE.getState();
default -> IotDeviceStateEnum.INACTIVE.getState();
};
}
}
private IotDeviceMessageServiceImpl getSelf() {
return SpringUtil.getBean(getClass());
}

View File

@@ -184,7 +184,8 @@ public class IotDevicePropertyServiceImpl implements IotDevicePropertyService {
processRuleProcessors(device, properties);
// 2.4 发布属性消息到 Redis Stream供其他模块如 Ops 订阅)
publishPropertyMessage(device, properties, message.getReportTime());
// TODO: 暂停发布,后续根据需要开启
// publishPropertyMessage(device, properties, message.getReportTime());
}
/**
@@ -197,7 +198,7 @@ public class IotDevicePropertyServiceImpl implements IotDevicePropertyService {
*/
private void processRuleProcessors(IotDeviceDO device, Map<String, Object> properties) {
try {
// 遍历所有属性,调用规<EFBFBD><EFBFBD>处理器
// 遍历所有属性,调用规处理器
for (Map.Entry<String, Object> entry : properties.entrySet()) {
String identifier = entry.getKey();
Object value = entry.getValue();
@@ -263,7 +264,7 @@ public class IotDevicePropertyServiceImpl implements IotDevicePropertyService {
.eventTime(reportTime)
.build();
integrationEventPublisher.publishPropertyChanged(event);
// integrationEventPublisher.publishPropertyChanged(event);
log.debug("[publishToIntegrationEventBus] 跨模块属性变更事件已发布: eventId={}, deviceId={}, productKey={}, properties={}",
event.getEventId(), device.getId(), productKey, properties.keySet());
} catch (Exception e) {