MqttService.java 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621
  1. package com.usky.sas.mqtt;
  2. import cn.hutool.core.bean.BeanUtil;
  3. import cn.hutool.core.util.IdUtil;
  4. import com.fasterxml.jackson.databind.DeserializationFeature;
  5. import com.fasterxml.jackson.databind.JsonNode;
  6. import com.fasterxml.jackson.databind.ObjectMapper;
  7. import com.usky.sas.domain.*;
  8. import com.usky.sas.mqtt.dto.*;
  9. import com.usky.sas.service.*;
  10. import lombok.extern.slf4j.Slf4j;
  11. import org.eclipse.paho.client.mqttv3.*;
  12. import org.springframework.beans.factory.annotation.Autowired;
  13. import org.springframework.beans.factory.annotation.Value;
  14. import org.springframework.stereotype.Service;
  15. import javax.annotation.PostConstruct;
  16. import javax.annotation.PreDestroy;
  17. import java.time.LocalDateTime;
  18. import java.time.format.DateTimeFormatter;
  19. import java.util.Arrays;
  20. import java.util.List;
  21. @Slf4j
  22. @Service
  23. public class MqttService {
  24. /** clientId、topics 仍从配置文件读取,连接参数 host/port/username/password 从 sas_config 表获取 */
  25. @Value("${mqtt.client-id:sas-client-${random.uuid}}")
  26. private String clientId;
  27. @Value("${mqtt.topics:}")
  28. private String topics;
  29. private MqttClient client;
  30. /** 是否处理消息:false 时收到消息不进入业务处理,用于启停监听 */
  31. private volatile boolean isListening = true;
  32. private final ObjectMapper mapper = new ObjectMapper()
  33. .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
  34. private final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
  35. @Autowired
  36. private SasUsbEventService usbEventService;
  37. @Autowired
  38. private SasSnapEventService snapEventService;
  39. @Autowired
  40. private SasEntranceEventService entranceEventService;
  41. @Autowired
  42. private SasParkingEventService parkingEventService;
  43. @Autowired
  44. private SasRoadblockEventService roadblockEventService;
  45. @Autowired
  46. private SasAlarsasEventService alarsasEventService;
  47. @Autowired
  48. private SasPatrolEventService patrolEventService;
  49. @Autowired
  50. private SasAcquisitionEventService acquisitionEventService;
  51. @Autowired
  52. private SasPerceptionEventService perceptionEventService;
  53. @Autowired
  54. private SasCollectionEventService collectionEventService;
  55. @Autowired
  56. private SasConfigService sasConfigService;
  57. @PostConstruct
  58. public void init() {
  59. SasConfig config = sasConfigService.getConfig();
  60. if (config == null || config.getHost() == null || config.getHost().isEmpty()) {
  61. log.warn("sas_config 中无 MQTT 配置或 host 为空,跳过 MQTT 连接");
  62. isListening = false;
  63. return;
  64. }
  65. try {
  66. String host = config.getHost();
  67. int port = parsePort(config.getPort(), 1883);
  68. String username = config.getUsername() != null ? config.getUsername() : "";
  69. String password = config.getPassword() != null ? config.getPassword() : "";
  70. String brokerUrl = "tcp://" + host + ":" + port;
  71. client = new MqttClient(brokerUrl, clientId);
  72. MqttConnectOptions options = new MqttConnectOptions();
  73. options.setCleanSession(true);
  74. if (!username.isEmpty()) {
  75. options.setUserName(username);
  76. options.setPassword(password.toCharArray());
  77. }
  78. options.setConnectionTimeout(60);
  79. options.setKeepAliveInterval(60);
  80. options.setAutomaticReconnect(true);
  81. client.setCallback(new MqttCallbackExtended() {
  82. @Override
  83. public void connectComplete(boolean reconnect, String serverURI) {
  84. log.info("MQTT连接成功, reconnect={}", reconnect);
  85. parseTopics().forEach(MqttService.this::subscribe);
  86. isListening = true;
  87. }
  88. @Override
  89. public void connectionLost(Throwable cause) {
  90. log.error("MQTT连接丢失: {}", cause.getMessage());
  91. isListening = false;
  92. scheduleReconnect();
  93. }
  94. @Override
  95. public void messageArrived(String topic, MqttMessage message) throws Exception {
  96. if (isListening) {
  97. processMessage(topic, message);
  98. }
  99. }
  100. @Override
  101. public void deliveryComplete(IMqttDeliveryToken token) {
  102. }
  103. });
  104. client.connect(options);
  105. log.info("MQTT客户端初始化完成, clientId={}", clientId);
  106. isListening = true;
  107. } catch (MqttException e) {
  108. log.error("MQTT初始化失败: {}", e.getMessage(), e);
  109. isListening = false;
  110. }
  111. }
  112. private static int parsePort(String portStr, int defaultPort) {
  113. if (portStr == null || portStr.isEmpty()) {
  114. return defaultPort;
  115. }
  116. try {
  117. return Integer.parseInt(portStr.trim());
  118. } catch (NumberFormatException e) {
  119. return defaultPort;
  120. }
  121. }
  122. /**
  123. * 连接丢失后延迟尝试主动重连(与自动重连互补)
  124. */
  125. private void scheduleReconnect() {
  126. Thread t = new Thread(() -> {
  127. try {
  128. Thread.sleep(5000L);
  129. if (client != null && !client.isConnected()) {
  130. log.info("尝试主动重连MQTT...");
  131. client.reconnect();
  132. }
  133. } catch (Exception e) {
  134. log.error("MQTT主动重连失败", e);
  135. }
  136. }, "mqtt-reconnect");
  137. t.setDaemon(true);
  138. t.start();
  139. }
  140. private List<String> parseTopics() {
  141. return Arrays.asList(topics.split(","));
  142. }
  143. private void processMessage(String topic, MqttMessage message) {
  144. String payload = new String(message.getPayload());
  145. log.debug("收到MQTT消息, topic={}, payload={}", topic, payload);
  146. try {
  147. // 从topic中提取事件类型,支持多级路径如 event/patrol/xxx 或 event/patrol
  148. String[] parts = topic.split("/");
  149. if (parts.length < 2) {
  150. log.warn("未知的topic格式: {}", topic);
  151. return;
  152. }
  153. // 查找事件类型关键字在topic中的位置
  154. String eventType = null;
  155. for (int i = 0; i < parts.length - 1; i++) {
  156. if ("event".equals(parts[i]) && i + 1 < parts.length) {
  157. eventType = parts[i + 1];
  158. break;
  159. }
  160. }
  161. if (eventType == null) {
  162. log.warn("无法从topic解析事件类型: {}", topic);
  163. return;
  164. }
  165. switch (eventType) {
  166. case "usbalarm":
  167. handleUsbEvent(payload, topic);
  168. break;
  169. case "snap":
  170. handleSnapEvent(payload, topic);
  171. break;
  172. case "entrance":
  173. handleEntranceEvent(payload, topic);
  174. break;
  175. case "parking":
  176. handleParkingEvent(payload, topic);
  177. break;
  178. case "roadblock":
  179. handleRoadblockEvent(payload, topic);
  180. break;
  181. case "alarm":
  182. handleAlarsasEvent(payload, topic);
  183. break;
  184. case "patrol":
  185. handlePatrolEvent(payload, topic);
  186. break;
  187. case "acquisition":
  188. handleAcquisitionEvent(payload, topic);
  189. break;
  190. case "perception":
  191. handlePerceptionEvent(payload, topic);
  192. break;
  193. case "collection":
  194. handleCollectionEvent(payload, topic);
  195. break;
  196. default:
  197. log.warn("未知的事件类型: {}, topic={}", eventType, topic);
  198. }
  199. } catch (Exception e) {
  200. log.error("处理MQTT消息失败, topic={}", topic, e);
  201. }
  202. }
  203. private void handleUsbEvent(String payload, String topic) {
  204. try {
  205. JsonNode root = mapper.readTree(payload);
  206. String method = root.path("method").asText();
  207. if ("event".equals(method)) {
  208. UsbEventMessage message = mapper.readValue(payload, UsbEventMessage.class);
  209. SasUsbEvent event = new SasUsbEvent();
  210. String eventId = message.getId() != null && !message.getId().isEmpty()
  211. ? message.getId()
  212. : IdUtil.randomUUID();
  213. event.setEventId(eventId);
  214. BeanUtil.copyProperties(message, event, "eventId");
  215. event.setCardId(message.getIcId());
  216. event.setTriggerTime(parseTime(message.getTriggerTime()));
  217. event.setCreateTime(LocalDateTime.now());
  218. event.setUpdateTime(LocalDateTime.now());
  219. usbEventService.save(event);
  220. log.info("USB事件保存成功, eventId={}", event.getEventId());
  221. } else if ("heart".equals(method)) {
  222. log.debug("收到USB设备心跳, deviceId={}", extractDeviceId(topic));
  223. }
  224. } catch (Exception e) {
  225. log.error("处理USB事件失败: {}", e.getMessage(), e);
  226. }
  227. }
  228. private void handleSnapEvent(String payload, String topic) {
  229. try {
  230. JsonNode root = mapper.readTree(payload);
  231. String method = root.path("method").asText();
  232. if ("event".equals(method)) {
  233. SnapEventMessage message = mapper.readValue(payload, SnapEventMessage.class);
  234. SasSnapEvent event = new SasSnapEvent();
  235. String eventId = message.getId() != null && !message.getId().isEmpty()
  236. ? message.getId()
  237. : IdUtil.randomUUID();
  238. // 如果该事件已存在(设备重发/重复消息),避免主键冲突,直接跳过
  239. if (snapEventService.getById(eventId) != null) {
  240. log.info("抓拍事件已存在, 跳过保存, eventId={}", eventId);
  241. return;
  242. }
  243. event.setEventId(eventId);
  244. BeanUtil.copyProperties(message, event, "eventId");
  245. event.setTriggerTime(parseTime(message.getTriggerTime()));
  246. event.setCreateTime(LocalDateTime.now());
  247. event.setUpdateTime(LocalDateTime.now());
  248. snapEventService.save(event);
  249. log.info("抓拍事件保存成功, eventId={}", event.getEventId());
  250. } else if ("heart".equals(method)) {
  251. log.debug("收到抓拍设备心跳, deviceId={}", extractDeviceId(topic));
  252. }
  253. } catch (Exception e) {
  254. log.error("处理抓拍事件失败: {}", e.getMessage(), e);
  255. }
  256. }
  257. private void handleEntranceEvent(String payload, String topic) {
  258. try {
  259. JsonNode root = mapper.readTree(payload);
  260. String method = root.path("method").asText();
  261. if ("event".equals(method)) {
  262. EntranceEventMessage message = mapper.readValue(payload, EntranceEventMessage.class);
  263. SasEntranceEvent event = new SasEntranceEvent();
  264. String eventId = message.getId() != null && !message.getId().isEmpty()
  265. ? message.getId()
  266. : IdUtil.randomUUID();
  267. event.setEventId(eventId);
  268. BeanUtil.copyProperties(message, event, "eventId");
  269. event.setCardId(message.getIcId());
  270. event.setTriggerTime(parseTime(message.getTriggerTime()));
  271. event.setCreateTime(LocalDateTime.now());
  272. event.setUpdateTime(LocalDateTime.now());
  273. entranceEventService.save(event);
  274. log.info("门禁事件保存成功, eventId={}", event.getEventId());
  275. } else if ("heart".equals(method)) {
  276. log.debug("收到门禁设备心跳, deviceId={}", extractDeviceId(topic));
  277. }
  278. } catch (Exception e) {
  279. log.error("处理门禁事件失败: {}", e.getMessage(), e);
  280. }
  281. }
  282. private void handleParkingEvent(String payload, String topic) {
  283. try {
  284. JsonNode root = mapper.readTree(payload);
  285. String method = root.path("method").asText();
  286. if ("event".equals(method)) {
  287. ParkingEventMessage message = mapper.readValue(payload, ParkingEventMessage.class);
  288. SasParkingEvent event = new SasParkingEvent();
  289. String eventId = message.getId() != null && !message.getId().isEmpty()
  290. ? message.getId()
  291. : IdUtil.randomUUID();
  292. event.setEventId(eventId);
  293. BeanUtil.copyProperties(message, event, "eventId");
  294. event.setTriggerTime(parseTime(message.getTriggerTime()));
  295. event.setCreateTime(LocalDateTime.now());
  296. event.setUpdateTime(LocalDateTime.now());
  297. parkingEventService.save(event);
  298. log.info("停车场事件保存成功, eventId={}", event.getEventId());
  299. } else if ("heart".equals(method)) {
  300. log.debug("收到停车场设备心跳, deviceId={}", extractDeviceId(topic));
  301. }
  302. } catch (Exception e) {
  303. log.error("处理停车场事件失败: {}", e.getMessage(), e);
  304. }
  305. }
  306. private void handleRoadblockEvent(String payload, String topic) {
  307. try {
  308. JsonNode root = mapper.readTree(payload);
  309. String method = root.path("method").asText();
  310. if ("event".equals(method)) {
  311. RoadblockEventMessage message = mapper.readValue(payload, RoadblockEventMessage.class);
  312. SasRoadblockEvent event = new SasRoadblockEvent();
  313. String eventId = message.getId() != null && !message.getId().isEmpty()
  314. ? message.getId()
  315. : IdUtil.randomUUID();
  316. event.setEventId(eventId);
  317. BeanUtil.copyProperties(message, event, "eventId");
  318. event.setTriggerTime(parseTime(message.getTriggerTime()));
  319. event.setCreateTime(LocalDateTime.now());
  320. event.setUpdateTime(LocalDateTime.now());
  321. roadblockEventService.save(event);
  322. log.info("阻车路障事件保存成功, eventId={}", event.getEventId());
  323. } else if ("heart".equals(method)) {
  324. log.debug("收到阻车路障设备心跳, deviceId={}", extractDeviceId(topic));
  325. }
  326. } catch (Exception e) {
  327. log.error("处理阻车路障事件失败: {}", e.getMessage(), e);
  328. }
  329. }
  330. private void handleAlarsasEvent(String payload, String topic) {
  331. try {
  332. JsonNode root = mapper.readTree(payload);
  333. String method = root.path("method").asText();
  334. if ("event".equals(method)) {
  335. AlarsasEventMessage message = mapper.readValue(payload, AlarsasEventMessage.class);
  336. SasAlarsasEvent event = new SasAlarsasEvent();
  337. String eventId = message.getId() != null && !message.getId().isEmpty()
  338. ? message.getId()
  339. : IdUtil.randomUUID();
  340. event.setEventId(eventId);
  341. BeanUtil.copyProperties(message, event, "eventId");
  342. event.setTriggerTime(parseTime(message.getTriggerTime()));
  343. event.setCreateTime(LocalDateTime.now());
  344. event.setUpdateTime(LocalDateTime.now());
  345. alarsasEventService.save(event);
  346. log.info("入侵报警事件保存成功, eventId={}", event.getEventId());
  347. } else if ("heart".equals(method)) {
  348. log.debug("收到入侵报警设备心跳, deviceId={}", extractDeviceId(topic));
  349. }
  350. } catch (Exception e) {
  351. log.error("处理入侵报警事件失败: {}", e.getMessage(), e);
  352. }
  353. }
  354. private void handlePatrolEvent(String payload, String topic) {
  355. try {
  356. JsonNode root = mapper.readTree(payload);
  357. String method = root.path("method").asText();
  358. if ("event".equals(method)) {
  359. PatrolEventMessage message = mapper.readValue(payload, PatrolEventMessage.class);
  360. SasPatrolEvent event = new SasPatrolEvent();
  361. String eventId = message.getId() != null && !message.getId().isEmpty()
  362. ? message.getId()
  363. : IdUtil.randomUUID();
  364. event.setEventId(eventId);
  365. BeanUtil.copyProperties(message, event, "eventId");
  366. event.setTriggerTime(parseTime(message.getTriggerTime()));
  367. event.setCreateTime(LocalDateTime.now());
  368. event.setUpdateTime(LocalDateTime.now());
  369. patrolEventService.save(event);
  370. log.info("巡检事件保存成功, eventId={}", event.getEventId());
  371. } else if ("heart".equals(method)) {
  372. log.debug("收到巡检设备心跳, deviceId={}", extractDeviceId(topic));
  373. }
  374. } catch (Exception e) {
  375. log.error("处理巡检事件失败: {}", e.getMessage(), e);
  376. }
  377. }
  378. private void handleAcquisitionEvent(String payload, String topic) {
  379. try {
  380. JsonNode root = mapper.readTree(payload);
  381. String method = root.path("method").asText();
  382. if ("event".equals(method)) {
  383. AcquisitionEventMessage message = mapper.readValue(payload, AcquisitionEventMessage.class);
  384. SasAcquisitionEvent event = new SasAcquisitionEvent();
  385. String eventId = message.getId() != null && !message.getId().isEmpty()
  386. ? message.getId()
  387. : IdUtil.randomUUID();
  388. event.setEventId(eventId);
  389. BeanUtil.copyProperties(message, event, "eventId");
  390. event.setTriggerTime(parseTime(message.getTriggerTime()));
  391. event.setCreateTime(LocalDateTime.now());
  392. event.setUpdateTime(LocalDateTime.now());
  393. acquisitionEventService.save(event);
  394. log.info("数据采集事件保存成功, eventId={}", event.getEventId());
  395. } else if ("heart".equals(method)) {
  396. log.debug("收到数据采集设备心跳, deviceId={}", extractDeviceId(topic));
  397. }
  398. } catch (Exception e) {
  399. log.error("处理数据采集事件失败", e);
  400. }
  401. }
  402. private void handlePerceptionEvent(String payload, String topic) {
  403. try {
  404. JsonNode root = mapper.readTree(payload);
  405. String method = root.path("method").asText();
  406. if ("event".equals(method)) {
  407. PerceptionEventMessage message = mapper.readValue(payload, PerceptionEventMessage.class);
  408. SasPerceptionEvent event = new SasPerceptionEvent();
  409. String eventId = message.getId() != null && !message.getId().isEmpty()
  410. ? message.getId()
  411. : IdUtil.randomUUID();
  412. event.setEventId(eventId);
  413. BeanUtil.copyProperties(message, event, "eventId");
  414. event.setTriggerTime(parseTime(message.getTriggerTime()));
  415. event.setCreateTime(LocalDateTime.now());
  416. event.setUpdateTime(LocalDateTime.now());
  417. perceptionEventService.save(event);
  418. log.info("状态感知事件保存成功, eventId={}", event.getEventId());
  419. } else if ("heart".equals(method)) {
  420. log.debug("收到状态感知设备心跳, deviceId={}", extractDeviceId(topic));
  421. }
  422. } catch (Exception e) {
  423. log.error("处理状态感知事件失败", e);
  424. }
  425. }
  426. private void handleCollectionEvent(String payload, String topic) {
  427. try {
  428. JsonNode root = mapper.readTree(payload);
  429. String method = root.path("method").asText();
  430. if ("event".equals(method)) {
  431. CollectionEventMessage message = mapper.readValue(payload, CollectionEventMessage.class);
  432. SasCollectionEvent event = new SasCollectionEvent();
  433. String eventId = message.getId() != null && !message.getId().isEmpty()
  434. ? message.getId()
  435. : IdUtil.randomUUID();
  436. event.setEventId(eventId);
  437. BeanUtil.copyProperties(message, event, "eventId");
  438. event.setTriggerTime(parseTime(message.getTriggerTime()));
  439. event.setCreateTime(LocalDateTime.now());
  440. event.setUpdateTime(LocalDateTime.now());
  441. collectionEventService.save(event);
  442. log.info("状态采集事件保存成功, eventId={}", event.getEventId());
  443. } else if ("heart".equals(method)) {
  444. log.debug("收到状态采集设备心跳, deviceId={}", extractDeviceId(topic));
  445. }
  446. } catch (Exception e) {
  447. log.error("处理状态采集事件失败", e);
  448. }
  449. }
  450. private String extractDeviceId(String topic) {
  451. String[] parts = topic.split("/");
  452. return parts.length > 0 ? parts[parts.length - 1] : "unknown";
  453. }
  454. private LocalDateTime parseTime(String timeStr) {
  455. if (timeStr == null || timeStr.isEmpty()) {
  456. return LocalDateTime.now();
  457. }
  458. try {
  459. return LocalDateTime.parse(timeStr, formatter);
  460. } catch (Exception e) {
  461. log.warn("时间解析失败: {}, 使用当前时间", timeStr);
  462. return LocalDateTime.now();
  463. }
  464. }
  465. public void subscribe(String topic) {
  466. try {
  467. if (client != null && client.isConnected()) {
  468. client.subscribe(topic);
  469. log.info("订阅主题成功: {}", topic);
  470. }
  471. } catch (MqttException e) {
  472. log.error("订阅主题失败: {}", topic, e);
  473. }
  474. }
  475. /**
  476. * 暂停监听:不再处理新到达的 MQTT 消息,连接保持
  477. */
  478. public void pauseListening() {
  479. this.isListening = false;
  480. log.info("MQTT监听已暂停");
  481. }
  482. /**
  483. * 恢复监听:继续处理 MQTT 消息
  484. */
  485. public void resumeListening() {
  486. this.isListening = true;
  487. log.info("MQTT监听已恢复");
  488. }
  489. /**
  490. * 查询连接状态:是否已连接 Broker
  491. */
  492. public boolean getConnectionStatus() {
  493. return client != null && client.isConnected();
  494. }
  495. /**
  496. * 查询监听状态:是否正在处理消息(未暂停)
  497. */
  498. public boolean getListeningStatus() {
  499. return isListening;
  500. }
  501. /**
  502. * 主动重连:先断开并关闭旧连接,再重新初始化连接并订阅
  503. */
  504. public void reconnect() {
  505. try {
  506. if (client != null) {
  507. if (client.isConnected()) {
  508. client.disconnect();
  509. }
  510. client.close();
  511. client = null;
  512. }
  513. isListening = true;
  514. init();
  515. } catch (Exception e) {
  516. log.error("MQTT重连失败", e);
  517. isListening = false;
  518. }
  519. }
  520. @PreDestroy
  521. public void disconnect() {
  522. try {
  523. if (client != null) {
  524. if (client.isConnected()) {
  525. client.disconnect();
  526. log.info("MQTT连接已断开");
  527. }
  528. client.close();
  529. client = null;
  530. }
  531. } catch (MqttException e) {
  532. log.error("断开MQTT连接失败", e);
  533. }
  534. }
  535. }