PublishClient.java 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362
  1. package com.ozs.service.utils;
  2. import com.alibaba.fastjson2.JSON;
  3. import com.alibaba.fastjson2.JSONObject;
  4. import com.alibaba.fastjson2.filter.Filter;
  5. import com.alibaba.fastjson2.filter.SimplePropertyPreFilter;
  6. import com.ozs.common.utils.StringUtils;
  7. import com.ozs.common.utils.sign.Md5Utils;
  8. import com.ozs.common.utils.stateSecrets.SM4Utils;
  9. import com.ozs.common.utils.uuid.IdUtils;
  10. import com.ozs.service.entity.BaseCameraManagement;
  11. import com.ozs.service.entity.vo.*;
  12. import lombok.extern.slf4j.Slf4j;
  13. import org.eclipse.paho.client.mqttv3.MqttClient;
  14. import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
  15. import org.eclipse.paho.client.mqttv3.MqttDeliveryToken;
  16. import org.eclipse.paho.client.mqttv3.MqttException;
  17. import org.eclipse.paho.client.mqttv3.MqttMessage;
  18. import org.eclipse.paho.client.mqttv3.MqttPersistenceException;
  19. import org.eclipse.paho.client.mqttv3.MqttTopic;
  20. import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
  21. import org.springframework.beans.factory.annotation.Autowired;
  22. import org.springframework.beans.factory.annotation.Value;
  23. import org.springframework.stereotype.Component;
  24. import org.springframework.util.ObjectUtils;
  25. import javax.annotation.PostConstruct;
  26. import java.util.ArrayList;
  27. import java.util.HashMap;
  28. import java.util.Map;
  29. import java.util.UUID;
  30. /**
  31. * 发布客户端
  32. *
  33. * @author Administrator
  34. */
  35. @Slf4j
  36. @Component
  37. public class PublishClient {
  38. /**
  39. * mqtt服务器地址
  40. */
  41. public static String HOST;
  42. public static String USER_NAME;
  43. public static String PASS_WORD;
  44. private static String clientId;
  45. /**
  46. * 客户端实例
  47. */
  48. private static MqttClient client;
  49. @Autowired
  50. private RabbitMqConfig rabbitMqConfig;
  51. /**
  52. * 获取 MqttTopic
  53. */
  54. public static MqttTopic getMqttTopic(String topic) {
  55. return client.getTopic(topic);
  56. }
  57. /**
  58. * 发布
  59. *
  60. * @param topic
  61. * @param message
  62. * @throws MqttPersistenceException
  63. * @throws MqttException
  64. */
  65. public static void publish(MqttTopic topic, MqttMessage message) throws MqttPersistenceException,
  66. MqttException {
  67. MqttDeliveryToken token = topic.publish(message);
  68. token.waitForCompletion();
  69. log.info("message is published completely! " + token.isComplete());
  70. }
  71. public static void pull(String sign, String deviceSn) throws MqttException {
  72. /**
  73. * 发布客户端
  74. */
  75. Heartbeat test = new Heartbeat();
  76. test.setName("HeartResponse");
  77. Data data = new Data();
  78. data.setAlarm(0);
  79. data.setStream(0);
  80. test.setData(data);
  81. test.setCode(200);
  82. test.setSign(sign);
  83. String s = JSON.toJSONString(test);
  84. MqttMessage message = new MqttMessage();
  85. /**
  86. * 保证消息能到达一次
  87. */
  88. message.setQos(1);
  89. /**
  90. * 消息保留
  91. */
  92. // server.message.setRetained(false);
  93. /**
  94. * 消息内容
  95. */
  96. message.setPayload(s.getBytes());
  97. /**
  98. * 发布
  99. */
  100. publish(getMqttTopic("heart_" + deviceSn), message);
  101. }
  102. public static void main(String[] args) {
  103. JSONObject res = new JSONObject();
  104. res.put("Name", "ConfigRequest");
  105. Codec codec = new Codec();
  106. Venc venc = new Venc();
  107. venc.setFps(Double.valueOf(11));
  108. ArrayList<Venc> vencList = new ArrayList<>();
  109. vencList.add(venc);
  110. codec.setVenc(vencList);
  111. Map<String, Object> map = new HashMap<>();
  112. res.put("device_sn", "11");
  113. res.put("sign", "rate" + IdUtils.fastSimpleUUID());
  114. map.put("codec",codec);
  115. res.put("data", map);
  116. String s = JSONObject.toJSONString(res);
  117. log.info(s);
  118. }
  119. public static void updateDeviceSn(BaseCameraVersionVo baseCameraVersionVo) {
  120. /**
  121. * 发布客户端
  122. */
  123. log.info("updateDeviceSn---start");
  124. for (String code : baseCameraVersionVo.getCameraCodeList()) {
  125. try {
  126. log.info("update_" + code);
  127. UpdateDeviceSn updateDeviceSn = new UpdateDeviceSn();
  128. updateDeviceSn.setName("UpdateRequest");
  129. updateDeviceSn.setType(Integer.valueOf(baseCameraVersionVo.getUpgradeType()));
  130. updateDeviceSn.setMd5(baseCameraVersionVo.getMd5());
  131. updateDeviceSn.setSign(IdUtils.fastSimpleUUID());
  132. updateDeviceSn.setUrl(baseCameraVersionVo.getVersionAddress());
  133. String s = JSON.toJSONString(updateDeviceSn);
  134. MqttMessage message = new MqttMessage();
  135. /**
  136. * 保证消息能到达一次
  137. */
  138. message.setQos(1);
  139. /**
  140. * 消息保留
  141. */
  142. // server.message.setRetained(false);
  143. /**
  144. * 消息内容
  145. */
  146. message.setPayload(s.getBytes());
  147. /**
  148. * 发布
  149. */
  150. publish(getMqttTopic("update_" + code), message);
  151. log.info("updateDeviceSn---end");
  152. } catch (MqttException e) {
  153. log.error("updateDeviceSn-------" + e.getMessage());
  154. }
  155. }
  156. }
  157. public static void alarmPush(ReqMsgAlarmVo reqMsgAlarmVo) {
  158. /**
  159. * 发布客户端
  160. */
  161. log.info("alarmPush---start");
  162. try {
  163. if (StringUtils.isEmpty(reqMsgAlarmVo.getAlarmMile())){
  164. reqMsgAlarmVo.setAlarmMile("");
  165. }
  166. if (ObjectUtils.isEmpty(reqMsgAlarmVo.getLineDir())){
  167. reqMsgAlarmVo.setLineDirStr("");
  168. }
  169. String s = JSON.toJSONString(reqMsgAlarmVo);
  170. MqttMessage message = new MqttMessage();
  171. /**
  172. * 保证消息能到达一次
  173. */
  174. message.setQos(1);
  175. /**
  176. * 消息保留
  177. */
  178. // server.message.setRetained(false);
  179. /**
  180. * 消息内容
  181. */
  182. message.setPayload(s.getBytes());
  183. /**
  184. * 发布
  185. */
  186. publish(getMqttTopic("alarmPush"), message);
  187. log.info("alarmPush---end");
  188. } catch (MqttException e) {
  189. log.error("alarmPush-------" + e.getMessage());
  190. }
  191. }
  192. public static void confidenceCoefficient(BaseCameraManagement baseCameraManagement, String value) {
  193. /**
  194. * 发布客户端
  195. */
  196. try {
  197. JSONObject res = new JSONObject();
  198. res.put("Name", "HeartRequest");
  199. Codec codec = new Codec();
  200. Venc venc = new Venc();
  201. venc.setFps(Double.valueOf(value));
  202. ArrayList<Venc> vencList = new ArrayList<>();
  203. vencList.add(venc);
  204. codec.setVenc(vencList);
  205. Map<String, Object> map = new HashMap<>();
  206. res.put("device_sn", baseCameraManagement.getCameraSn());
  207. res.put("sign", "rate" + IdUtils.fastSimpleUUID());
  208. map.put("codec",codec);
  209. res.put("data", map);
  210. String s = JSONObject.toJSONString(res);
  211. MqttMessage message = new MqttMessage();
  212. /**
  213. * 保证消息能到达一次
  214. */
  215. message.setQos(1);
  216. /**
  217. * 消息保留
  218. */
  219. // message.setRetained(false);
  220. /**
  221. * 消息内容
  222. */
  223. message.setPayload(s.getBytes());
  224. /**
  225. * 发布
  226. */
  227. publish(getMqttTopic("config_" + baseCameraManagement.getCameraCode()), message);
  228. } catch (MqttException e) {
  229. log.error(e.getMessage());
  230. }
  231. }
  232. public static void configFrameRate(BaseCameraManagement baseCameraManagement,Integer mode) {
  233. /**
  234. * 发布客户端
  235. */
  236. try {
  237. JSONObject res = new JSONObject();
  238. res.put("Name", "HeartRequest");
  239. Codec codec = new Codec();
  240. Day2night day2night = new Day2night();
  241. day2night.setMode(mode);
  242. codec.setDay2night(day2night);
  243. Map<String, Object> map = new HashMap<>();
  244. res.put("device_sn", baseCameraManagement.getCameraSn());
  245. res.put("sign", "rate" + IdUtils.fastSimpleUUID());
  246. map.put("codec",codec);
  247. res.put("data", map);
  248. String s = JSONObject.toJSONString(res);
  249. MqttMessage message = new MqttMessage();
  250. /**
  251. * 保证消息能到达一次
  252. */
  253. message.setQos(1);
  254. /**
  255. * 消息保留
  256. */
  257. // message.setRetained(false);
  258. /**
  259. * 消息内容
  260. */
  261. message.setPayload(s.getBytes());
  262. /**
  263. * 发布
  264. */
  265. publish(getMqttTopic("config_" + baseCameraManagement.getCameraCode()), message);
  266. } catch (MqttException e) {
  267. log.error(e.getMessage());
  268. }
  269. }
  270. @PostConstruct
  271. public void init() {
  272. HOST = rabbitMqConfig.getHost();
  273. USER_NAME = rabbitMqConfig.getUserName();
  274. PASS_WORD = rabbitMqConfig.getPassword();
  275. clientId = rabbitMqConfig.getClientId();
  276. /**
  277. * 连接配置
  278. */
  279. MqttConnectOptions options = new MqttConnectOptions();
  280. /**
  281. * 不保存,每次重启新client
  282. */
  283. options.setCleanSession(true);
  284. options.setUserName(USER_NAME);
  285. options.setPassword(PASS_WORD.toCharArray());
  286. /**
  287. * 设置超时时间
  288. */
  289. options.setConnectionTimeout(60);
  290. /**
  291. * 设置会话心跳时间
  292. */
  293. options.setKeepAliveInterval(40);
  294. try {
  295. /**
  296. * 设置发布回调
  297. */
  298. client = new MqttClient(HOST, clientId, new MemoryPersistence());
  299. client.setCallback(new PublishCallback());
  300. client.connect(options);
  301. int[] Qos = {1};
  302. String[] topic1 = {"config", "update", "heart", "test"};
  303. client.subscribe(topic1);
  304. } catch (Exception e) {
  305. log.error(e.getMessage());
  306. e.printStackTrace();
  307. }
  308. }
  309. public static void reconnect() throws MqttException {
  310. log.error("尝试重连...");
  311. if (client != null) {
  312. try {
  313. client.close();
  314. } catch (MqttException e) {
  315. log.error("关闭现有连接时出错:" + e);
  316. }
  317. }
  318. MqttConnectOptions options = new MqttConnectOptions();
  319. options.setCleanSession(true);
  320. options.setUserName("camera-update");
  321. options.setPassword("05J5+mtYzx.Ry".toCharArray());
  322. options.setConnectionTimeout(60);
  323. options.setKeepAliveInterval(40);
  324. try {
  325. /**
  326. * 设置发布回调
  327. */
  328. client = new MqttClient("tcp://10.161.12.60:1883", "HAZARD-CAMERA-CLIENTID-123", new MemoryPersistence());
  329. client.setCallback(new PublishCallback());
  330. client.connect(options);
  331. String[] topic1 = {"config", "update", "heart", "test"};
  332. client.subscribe(topic1);
  333. } catch (Exception e) {
  334. log.error(e.getMessage());
  335. e.printStackTrace();
  336. }
  337. log.error("重连成功!");
  338. }
  339. }