#更新优化

This commit is contained in:
2026-09-02 05:03:10 +00:00
parent a92ec91cd2
commit 5f369a159d
3 changed files with 83 additions and 51 deletions

View File

@@ -367,7 +367,7 @@ public class NettyDevice extends Connector {
}, 0, 5, TimeUnit.SECONDS);
}
} catch (Exception e) {
throw new RuntimeException(e);
log.error("发送开始任务:error {} ", e.getMessage());
}
}

View File

@@ -43,4 +43,6 @@ public class SysSite extends BaseDO {
private boolean dingTalkOrderPush = false;// 工单是否推送
private String siteMapName; // geoserver 地图名称
}

View File

@@ -20,7 +20,6 @@ import org.springframework.util.CollectionUtils;
import java.util.ArrayList;
import java.util.List;
/**
* mqtt相关的 上位机处理器
*/
@@ -33,7 +32,6 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
private final WebsocketMesDispather websocketMesDispather = SpringUtils.getBean(WebsocketMesDispather.class);
private final String mqttClientId = "maibu-mqtt-client";
@Override
@@ -49,8 +47,10 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
/**
* 处理 MQTT 消息
*
* <p>通配符订阅场景下,deviceKey 为配置中的固定标识(如 "host-monitor"),
* 真正的设备 ID 需从 topic 中解析,通过 {@link #extractDeviceId(String)} 获取。</p>
* <p>
* 通配符订阅场景下,deviceKey 为配置中的固定标识(如 "host-monitor"),
* 真正的设备 ID 需从 topic 中解析,通过 {@link #extractDeviceId(String)} 获取。
* </p>
*/
@Override
public void handleMessage(String deviceKey, String topic, String payload, CustomMqttDeviceConfig config) {
@@ -82,8 +82,8 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
break;
// case "route_post":
// case "point_post":
// handleTaskPost(deviceId, topic, payload, action);
// break;
// handleTaskPost(deviceId, topic, payload, action);
// break;
case "task_status_post":
handleTaskPost(deviceId, topic, payload, action);
break;
@@ -121,21 +121,22 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
try {
websocketMesDispather.dispather(x.getConnectorId(), JsonUtils.toJsonString(dto));
} catch (Exception e) {
log.error("位置数据websocket推送失败: connectorId={}, error={}", x.getConnectorId(), e.getMessage());
log.error("位置数据websocket推送失败: connectorId={}, error={}", x.getConnectorId(),
e.getMessage());
}
}
}
//todo 推送实时位置和状态消息到外部
//todo 记录到内存 开线程池队列 逐一转发
// GlobalMemory.customMqttDeviceMonitor.publish(mqttClientId, String.format(MqttTopic.DEVICE_LOCATION_TOPIC, deviceId),
// JsonUtils.toJsonString(dto));
// todo 推送实时位置和状态消息到外部
// todo 记录到内存 开线程池队列 逐一转发
// GlobalMemory.customMqttDeviceMonitor.publish(mqttClientId,
// String.format(MqttTopic.DEVICE_LOCATION_TOPIC, deviceId),
// JsonUtils.toJsonString(dto));
}
} catch (Exception e) {
log.error("位置消息解析失败: deviceId={}, payload={}, error={}", deviceId, payload, e.getMessage());
}
}
public WebStatusMessageDTO createWebDeviceLocationMessage(HostLocationRelTimeDTO locationRelTimeDTO) {
WebStatusMessageDTO webStatusMessageDTO = new WebStatusMessageDTO();
List<DeviceStatusDetail> transferData = new ArrayList<>();
@@ -152,13 +153,31 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
lng.setUnit("");
transferData.add(lng);
DeviceStatusDetail yaw = new DeviceStatusDetail();
yaw.setName("yaw");
yaw.setValue(String.valueOf(locationRelTimeDTO.getData().getYaw()));
yaw.setUnit("");
transferData.add(yaw);
DeviceStatusDetail fix_status = new DeviceStatusDetail();
fix_status.setName("fix_status");
fix_status.setValue(String.valueOf(locationRelTimeDTO.getData().getFix_status()));
fix_status.setUnit("");
transferData.add(fix_status);
DeviceStatusDetail altitude = new DeviceStatusDetail();
altitude.setName("altitude");
altitude.setValue(String.valueOf(locationRelTimeDTO.getData().getAltitude()));
altitude.setUnit("");
transferData.add(altitude);
webStatusMessageDTO.setData(transferData);
webStatusMessageDTO.setType(MesType.locationInfo);
return webStatusMessageDTO;
}
public WebStatusMessageDTO createWebDeviceRealtimeMessage(HostVehicleRelTimeDTO vehicleRelTimeDTO) {
//todo 状态消息
// todo 状态消息
WebStatusMessageDTO webStatusMessageDTO = new WebStatusMessageDTO();
List<DeviceStatusDetail> transferData = new ArrayList<>();
@@ -208,12 +227,14 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
try {
websocketMesDispather.dispather(x.getConnectorId(), JsonUtils.toJsonString(dto));
} catch (Exception e) {
log.error("实时数据websocket推送失败: connectorId={}, error={}", x.getConnectorId(), e.getMessage());
log.error("实时数据websocket推送失败: connectorId={}, error={}", x.getConnectorId(),
e.getMessage());
}
}
}
// GlobalMemory.customMqttDeviceMonitor.publish(mqttClientId, String.format(MqttTopic.DEVICE_STATUS_TOPIC, deviceId),
// JsonUtils.toJsonString(dto));
// GlobalMemory.customMqttDeviceMonitor.publish(mqttClientId,
// String.format(MqttTopic.DEVICE_STATUS_TOPIC, deviceId),
// JsonUtils.toJsonString(dto));
}
} catch (Exception e) {
@@ -257,7 +278,6 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
private void handleTaskPost(String deviceId, String topic, String payload, String action) {
log.info("【任务上报】deviceId={}, action={}, payload={}", deviceId, action, payload);
HostTaskStatusDTO dto = JsonUtils.parseObject(payload, HostTaskStatusDTO.class);
if (dto != null && !StringUtils.isEmpty(dto.getStatus())) {
NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
@@ -267,25 +287,24 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
if ("cancelled".equals(dto.getStatus())) {
nettyDevice.canceledTask();
}
// if("paused".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("failed".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("accepted".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("running".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("paused".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("failed".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("accepted".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
// if("running".equals(dto.getStatus())){
// NettyDevice nettyDevice = deviceSessionManager.getDevice(deviceId);
// nettyDevice.finishTask();
// }
}
}
/**
@@ -302,24 +321,23 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
log.debug("【控制响应】deviceId={}, payload={}", deviceId, payload);
}
/**
* 处理控制指令响应
*/
private void handleRouteReply(String deviceId, String topic, String payload) {
log.debug("【路径任务下发响应】deviceId={}, payload={}", deviceId, payload);
//todo 更新任务接收状态
// todo 更新任务接收状态
HostNavigationReplyDTO replyDTO = JsonUtils.parseObject(payload, HostNavigationReplyDTO.class);
NettyDevice device = deviceSessionManager.getDevice(deviceId);
if (replyDTO.isAccepted()) {
if (device != null && device.getTask() != null && String.valueOf(device.getTask().getId()).equals(replyDTO.getTask_id())) {
if (device != null && device.getTask() != null
&& String.valueOf(device.getTask().getId()).equals(replyDTO.getTask_id())) {
device.getTask().setTaskAccepted(true);
}
}
}
/**
* 从 Topic 中解析 deviceId
* Topic 格式:mower/{deviceId}/...
@@ -346,21 +364,30 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
}
StringBuilder path = new StringBuilder();
for (int i = 2; i < parts.length; i++) {
if (i > 2) path.append("/");
if (i > 2)
path.append("/");
path.append(parts[i]);
}
String remaining = path.toString();
if (remaining.equals("property/location/post")) return "location_post";
if (remaining.equals("property/realtime/post")) return "realtime_post";
if (remaining.equals("event/heartbeat/post")) return "heartbeat_post";
if (remaining.equals("event/error/post")) return "error_post";
if (remaining.equals("property/location/post"))
return "location_post";
if (remaining.equals("property/realtime/post"))
return "realtime_post";
if (remaining.equals("event/heartbeat/post"))
return "heartbeat_post";
if (remaining.equals("event/error/post"))
return "error_post";
// if (remaining.equals("task/route/post")) return "route_post";
// if (remaining.equals("task/point/post")) return "point_post";
if (remaining.equals("task/status/post")) return "task_status_post";
if (remaining.equals("action/config/reply")) return "config_reply";
if (remaining.equals("action/control/reply")) return "control_reply";
if (remaining.equals("task/target/reply")) return "route_reply";
if (remaining.equals("task/status/post"))
return "task_status_post";
if (remaining.equals("action/config/reply"))
return "config_reply";
if (remaining.equals("action/control/reply"))
return "control_reply";
if (remaining.equals("task/target/reply"))
return "route_reply";
return null;
}
@@ -369,7 +396,8 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
for (String name : names) {
JsonNode fn = node.get(name);
if (fn != null && !fn.isNull()) {
if (fn.isNumber()) return fn.asDouble();
if (fn.isNumber())
return fn.asDouble();
try {
return Double.parseDouble(fn.asText());
} catch (Exception ignored) {
@@ -383,7 +411,8 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
for (String name : names) {
JsonNode fn = node.get(name);
if (fn != null && !fn.isNull()) {
if (fn.isNumber()) return fn.asDouble();
if (fn.isNumber())
return fn.asDouble();
try {
return Double.parseDouble(fn.asText());
} catch (Exception ignored) {
@@ -397,7 +426,8 @@ public class HostMessageHandler implements CustomMqttMessageHandler {
for (String name : names) {
JsonNode fn = node.get(name);
if (fn != null && !fn.isNull()) {
if (fn.isNumber()) return fn.asInt();
if (fn.isNumber())
return fn.asInt();
try {
return Integer.parseInt(fn.asText());
} catch (Exception ignored) {