通配符订阅场景下,deviceKey 为配置中的固定标识(如 "host-monitor"), + * 真正的设备 ID 需从 topic 中解析,通过 {@link #extractDeviceId(String)} 获取。
+ */ + @Override + public void handleMessage(String deviceKey, String topic, String payload, CustomMqttDeviceConfig config) { + try { + String deviceId = extractDeviceId(topic); + if (deviceId == null) { + log.warn("无法从topic解析deviceId: {}", topic); + return; + } + + String action = extractAction(topic); + if (action == null) { + log.debug("未识别的topic: {}", topic); + return; + } + + switch (action) { + case "location_post": + handleLocationPost(deviceId, topic, payload); + break; + case "realtime_post": + handleRealtimePost(deviceId, topic, payload); + break; + case "status_post": + handleStatusPost(deviceId, topic, payload); + break; + case "heartbeat_post": + handleHeartbeatPost(deviceId, topic, payload); + break; + case "error_post": + handleErrorPost(deviceId, topic, payload); + break; + case "route_post": + case "point_post": + handleTaskPost(deviceId, topic, payload, action); + break; + case "task_status_post": + case "task_finish_post": + handleTaskPost(deviceId, topic, payload, action); + break; + case "config_reply": + handleConfigReply(deviceId, topic, payload); + break; + case "control_reply": + handleControlReply(deviceId, topic, payload); + break; + default: + log.debug("割草机消息: deviceId={}, topic={}, action={}", deviceId, topic, action); + break; + } + } catch (Exception e) { + log.error("消息处理异常: topic={}, error={}", topic, e.getMessage()); + } + } + + /** + * 处理 GPS 位置数据上报 + * topic: mower/{deviceId}/property/location/post + */ + private void handleLocationPost(String deviceId, String topic, String payload) { + try { + JsonNode jsonNode = objectMapper.readTree(payload); + + double latitude = getDouble(jsonNode, "latitude", "lat"); + double longitude = getDouble(jsonNode, "longitude", "lng", "lon"); + double speed = getDouble(jsonNode, "speed"); + double heading = getDouble(jsonNode, "heading", "direction"); + double altitude = getDouble(jsonNode, "altitude", "alt"); + int satellites = getInt(jsonNode, "satellites", "sats"); + String fixQuality = getString(jsonNode, "fixQuality", "fix_quality"); + Double batteryLevel = getDoubleObj(jsonNode, "batteryLevel", "battery"); + + log.info("【位置数据】deviceId={}, lat={}, lng={}, speed={}, heading={}, sats={}", + deviceId, latitude, longitude, speed, heading, satellites); + + onLocationUpdate(deviceId, latitude, longitude, speed, heading, altitude, satellites, fixQuality, batteryLevel); + + } catch (Exception e) { + log.error("位置消息解析失败: deviceId={}, payload={}, error={}", deviceId, payload, e.getMessage()); + } + } + + /** + * 处理实时运行数据 + * topic: mower/{deviceId}/property/realtime/post + */ + private void handleRealtimePost(String deviceId, String topic, String payload) { + log.debug("【实时数据】deviceId={}, payload={}", deviceId, payload); + } + + /** + * 处理状态变更事件 + * topic: mower/{deviceId}/event/status/post + */ + private void handleStatusPost(String deviceId, String topic, String payload) { + log.info("【状态变更】deviceId={}, payload={}", deviceId, payload); + } + + /** + * 处理心跳 + * topic: mower/{deviceId}/event/heartbeat/post + */ + private void handleHeartbeatPost(String deviceId, String topic, String payload) { + log.debug("【心跳】deviceId={}", deviceId); + } + + /** + * 处理错误告警 + * topic: mower/{deviceId}/event/error/post + */ + private void handleErrorPost(String deviceId, String topic, String payload) { + log.warn("【错误告警】deviceId={}, payload={}", deviceId, payload); + } + + /** + * 处理任务相关上报 + */ + private void handleTaskPost(String deviceId, String topic, String payload, String action) { + log.info("【任务上报】deviceId={}, action={}, payload={}", deviceId, action, payload); + } + + /** + * 处理配置响应 + */ + private void handleConfigReply(String deviceId, String topic, String payload) { + log.debug("【配置响应】deviceId={}, payload={}", deviceId, payload); + } + + /** + * 处理控制指令响应 + */ + private void handleControlReply(String deviceId, String topic, String payload) { + log.debug("【控制响应】deviceId={}, payload={}", deviceId, payload); + } + + /** + * 位置数据回调,子类可覆盖实现业务逻辑 + */ + protected void onLocationUpdate(String deviceId, double latitude, double longitude, + double speed, double heading, double altitude, + int satellites, String fixQuality, Double batteryLevel) { + } + + /** + * 从 Topic 中解析 deviceId + * Topic 格式:mower/{deviceId}/... + */ + private String extractDeviceId(String topic) { + if (topic == null || topic.isEmpty()) { + return null; + } + String[] parts = topic.split("/"); + if (parts.length < 2 || !"mower".equals(parts[0])) { + return null; + } + return parts[1]; + } + + /** + * 从 Topic 中提取动作标识 + * 将 mower/{deviceId}/xxx/yyy/post 转为 xxx_yyy_post 形式 + */ + private String extractAction(String topic) { + String[] parts = topic.split("/"); + if (parts.length < 3 || !"mower".equals(parts[0])) { + return null; + } + StringBuilder path = new StringBuilder(); + for (int i = 2; i < parts.length; i++) { + 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/status/post")) return "status_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("task/finish/post")) return "task_finish_post"; + if (remaining.equals("action/config/reply")) return "config_reply"; + if (remaining.equals("action/control/reply")) return "control_reply"; + + return null; + } + + private double getDouble(JsonNode node, String... names) { + for (String name : names) { + JsonNode fn = node.get(name); + if (fn != null && !fn.isNull()) { + if (fn.isNumber()) return fn.asDouble(); + try { + return Double.parseDouble(fn.asText()); + } catch (Exception ignored) { + } + } + } + return -1.0; + } + + private Double getDoubleObj(JsonNode node, String... names) { + for (String name : names) { + JsonNode fn = node.get(name); + if (fn != null && !fn.isNull()) { + if (fn.isNumber()) return fn.asDouble(); + try { + return Double.parseDouble(fn.asText()); + } catch (Exception ignored) { + } + } + } + return null; + } + + private int getInt(JsonNode node, String... names) { + for (String name : names) { + JsonNode fn = node.get(name); + if (fn != null && !fn.isNull()) { + if (fn.isNumber()) return fn.asInt(); + try { + return Integer.parseInt(fn.asText()); + } catch (Exception ignored) { + } + } + } + return -1; + } + + private String getString(JsonNode node, String... names) { + for (String name : names) { + JsonNode fn = node.get(name); + if (fn != null && !fn.isNull()) { + return fn.asText(); + } + } + return null; + } + + @Override + public void onConnectionLost(String deviceKey, Throwable cause) { + log.warn("主机监控连接断开: deviceKey={}, reason={}", deviceKey, cause.getMessage()); + } + + @Override + public void onSubscribed(String deviceKey, String topic) { + log.info("主机监控订阅主题: deviceKey={}, topic={}", deviceKey, topic); + } +} diff --git a/maibu-common/src/main/java/com/maibu/mqtt/MqttTopic.java b/maibu-common/src/main/java/com/maibu/mqtt/MqttTopic.java index d926eb1..7a9a90e 100644 --- a/maibu-common/src/main/java/com/maibu/mqtt/MqttTopic.java +++ b/maibu-common/src/main/java/com/maibu/mqtt/MqttTopic.java @@ -18,12 +18,12 @@ public class MqttTopic { // ==================== 设备连接与状态 ==================== /** 设备上线/下线事件上报(设备->平台),对应 TCP 0x03 connect */ public static final String MOWER_STATUS_POST = "mower/%s/event/status/post"; - /** 状态变更确认回复(平台->设备) */ -// public static final String MOWER_STATUS_REPLY = "mower/%s/event/status/reply"; // ==================== 实时运行数据 ==================== /** 实时运行数据上报:电压/速度/温度/位置等(设备->平台),对应 TCP 0x02 status */ public static final String MOWER_REALTIME_POST = "mower/%s/property/realtime/post"; + /** GPS 位置数据上报:经纬度、速度、航向等(设备->平台) */ + public static final String MOWER_LOCATION_POST = "mower/%s/property/location/post"; // ==================== 心跳保活 ==================== /** 设备心跳上报(设备->平台),对应 TCP 0xFF heartbeat */ @@ -43,22 +43,12 @@ public class MqttTopic { /** 远程控制指令下发:前进/后退/转向/割草等(平台->设备),对应 TCP 0x00 remoteControl */ public static final String MOWER_CONTROL_SET = "mower/%s/action/control/set"; -// /** 控制指令响应(设备->平台) */ -// public static final String MOWER_CONTROL_REPLY = "mower/%s/action/control/reply"; - // ==================== 任务管理 ==================== /** 路径规划下发(平台->设备),对应 TCP 0x01 path */ public static final String MOWER_TASK_ROUTE_SET = "mower/%s/task/route/set"; public static final String MOWER_TASK_ROUTE_POST = "mower/%s/task/route/post"; /** 启动任务执行(平台->设备),对应 TCP 0x12 interaction */ public static final String MOWER_TASK_STATUS_SET = "mower/%s/task/status/set"; -// /** 暂停任务(平台->设备) */ -// public static final String MOWER_TASK_PAUSE_SET = "mower/%s/task/pause/set"; -// /** 继续任务(平台->设备) */ -// public static final String MOWER_TASK_RESUME_SET = "mower/%s/task/resume/set"; -// /** 取消任务(平台->设备) */ -// public static final String MOWER_TASK_CANCEL_SET = "mower/%s/task/cancel/set"; -// /** 到达路径点上报(设备->平台) */ public static final String MOWER_TASK_POINT_POST = "mower/%s/task/point/post"; /** 任务状态变更上报:执行中/暂停/完成等(设备->平台) */ public static final String MOWER_TASK_STATUS_POST = "mower/%s/task/status/post"; @@ -90,6 +80,8 @@ public class MqttTopic { public static final String MOWER_WILDCARD_STATUS = "mower/+/event/status/post"; /** 订阅所有设备的实时运行数据 */ public static final String MOWER_WILDCARD_REALTIME = "mower/+/property/realtime/post"; + /** 上位机的额外消息 */ + public static final String MOWER_WILDCARD_LOCATION = "mower/+/property/location/post"; /** 订阅所有设备的错误告警 */ public static final String MOWER_WILDCARD_ERROR = "mower/+/event/error/post"; /** 订阅所有设备的任务状态变更 */ diff --git a/maibu-netty-server/src/main/java/com/maibu/init/InitThread.java b/maibu-netty-server/src/main/java/com/maibu/init/InitThread.java index 774a769..44be4c2 100644 --- a/maibu-netty-server/src/main/java/com/maibu/init/InitThread.java +++ b/maibu-netty-server/src/main/java/com/maibu/init/InitThread.java @@ -121,8 +121,9 @@ public class InitThread implements ApplicationRunner { private void init() throws MqttException { // 初始化InfluxDBUtil GlobalMemory.influxDBUtil = new InfluxDBUtil(influxDBUrl, influxDBToken, influxDBOrg, influxDBBucket); - // 初始化MqttClientUtil - GlobalMemory.mqttClientUtil = new MqttClientUtil(mqttHostUrl, mqttClientId, mqttUsername, mqttPassword); + + GlobalMemory.initMqttDeviceMonitor(mqttHostUrl, mqttClientId, mqttUsername, mqttPassword); + // 设备错误监听线程 deviceErrorMonitorExecutor.execute(() -> { try { diff --git a/maibu-netty-server/src/main/java/com/maibu/netty/handler/DataToDataBaseHandler.java b/maibu-netty-server/src/main/java/com/maibu/netty/handler/DataToDataBaseHandler.java index 87760a1..69b408d 100644 --- a/maibu-netty-server/src/main/java/com/maibu/netty/handler/DataToDataBaseHandler.java +++ b/maibu-netty-server/src/main/java/com/maibu/netty/handler/DataToDataBaseHandler.java @@ -20,6 +20,7 @@ import org.eclipse.paho.client.mqttv3.MqttException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.springframework.util.CollectionUtils; @@ -63,6 +64,10 @@ public class DataToDataBaseHandler extends SimpleChannelInboundHandler