ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

spring boot 连接emqx并实现发布订阅

spring boot 连接emqx并实现发布订阅 1、安装依赖dependency groupIdorg.springframework.integration/groupId artifactIdspring-integration-mqtt/artifactId version7.1.0/version scopecompile/scope /dependency !-- Source: https://mvnrepository.com/artifact/org.eclipse.paho/org.eclipse.paho.client.mqttv3 -- dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.2.5/version scopecompile/scope /dependency2、配置连接信息#MQTT???? #MQTT-??? spring.mqtt.usernamexxxxxx #MQTT-?? spring.mqtt.passwordxxxxx #MQTT-??????????????????????tcp://127.0.0.1:61613?tcp://47.123.33.66:61613 spring.mqtt.urltcp://127.0.0.1:1883 #MQTT-??????????ID spring.mqtt.client-idemq2 #MQTT-????????????????????? spring.mqtt.topicxiaomingming #timeout ?????? spring.mqtt.timeout20 #keep alive spring.mqtt.keep-alive20 spring.mqtt.qos23、新建mqtt配置类package com.example.springbootemq.config; import org.eclipse.paho.client.mqttv3.MqttClient; import org.eclipse.paho.client.mqttv3.MqttConnectOptions; import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Configuration; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; Component Configuration public class MqttConfig { Value(${spring.mqtt.username}) private String username; Value(${spring.mqtt.password}) private String password; Value(${spring.mqtt.url}) private String hostUrl; Value(${spring.mqtt.client-id}) private String clientId; Value(${spring.mqtt.timeout}) private Integer timeout; Value(${spring.mqtt.keep-alive}) private Integer keepAlive; Value(${spring.mqtt.qos}) private Integer qos; Value(${spring.mqtt.topic}) private String topic; private MqttClient client; /** * 项目启动时自动连接 MQTT */ PostConstruct public void init() { connect(); } //断线手动重连重新订阅 public void reConnectSubscribe() { while (true){ try { log.warn(开始重连重订阅); Thread.sleep(2000); this.client.connect(connOpts); //订阅 this.client.subscribe(myTest,2); break; }catch (Exception ex){ ex.printStackTrace(); } } } /** * 连接 MQTT */ public void connect() { try { client new MqttClient(hostUrl, clientId, new MemoryPersistence()); // MQTT 连接选项 MqttConnectOptions connOpts new MqttConnectOptions(); connOpts.setUserName(username); connOpts.setPassword(password.toCharArray()); // 保留会话 connOpts.setCleanSession(true); // 设置超时时间单位秒 connOpts.setConnectionTimeout(timeout); // 设置心跳时间单位秒表示服务器每隔1.5*20秒的时间向客户端发送心跳判断客户端是否在线 connOpts.setKeepAliveInterval(keepAlive); // 设置回调 client.setCallback(new OnMessageCallback()); // 建立连接 client.connect(connOpts); //订阅 client.subscribe(topic,2); } catch (MqttException me) { System.out.println(reason me.getReasonCode()); System.out.println(msg me.getMessage()); System.out.println(loc me.getLocalizedMessage()); System.out.println(cause me.getCause()); System.out.println(excep me); me.printStackTrace(); } } /** * 订阅 * * param topic 主题 */ public void subscribe(String topic) { try { client.subscribe(topic, qos); } catch (MqttException me) { me.printStackTrace(); } } /** * 消息发布 * * param topic 主题 * param data 消息 */ public void publish(String topic, String data) { try { MqttMessage message new MqttMessage(data.getBytes()); message.setQos(qos); // 消息服务质量等级 message.setRetained(true); // 保留消息 client.publish(topic, message); } catch (MqttException me) { me.printStackTrace(); } } /** * 断开连接 */ public void disconnect() { try { client.disconnect(); client.close(); } catch (MqttException me) { me.printStackTrace(); } } }4、定义消息回调类OnMessageCallbackpackage com.example.springbootemq.config; import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken; import org.eclipse.paho.client.mqttv3.MqttCallback; import org.eclipse.paho.client.mqttv3.MqttMessage; public class OnMessageCallback implements MqttCallback { private MqttConfig mqttConfig; //构造函数注入对象 public OnMessageCallback(MqttConfig mqttConfig) { this.mqttConfig mqttConfig; } Override public void connectionLost(Throwable cause) { // 连接丢失后一般在这里面进行重连 System.out.println(连接断开可以做重连); this.mqttConfig.reConnectSubscribe(); } Override public void messageArrived(String topic, MqttMessage message) throws Exception { // subscribe后得到的消息会执行到这里面 System.out.println(接收消息主题: topic); System.out.println(接收消息Qos: message.getQos()); System.out.println(接收消息内容: new String(message.getPayload())); } Override public void deliveryComplete(IMqttDeliveryToken token) { System.out.println(deliveryComplete--------- token.isComplete()); } }5、控制器测试package com.example.springbootemq.controller; import com.example.springbootemq.config.MqttConfig; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; RestController RequestMapping(/mqtt) public class MqttController { Resource private MqttConfig mqttConfig; GetMapping(/send/{message}) public void sendMessage(PathVariable(message) String message) { String subTopic testtopic/#; String pubTopic testtopic/1; String topicTestzmtest; String data hello MQTT test; mqttConfig.publish(topicTest, message); // // 订阅 // mqttConfig.subscribe(subTopic); // // 发布消息 // mqttConfig.publish(pubTopic, data); // 断开连接 // mqttConfig.disconnect(); } GetMapping(/receive) public void receiveMessage() { String subTopic testtopic/#; String pubTopic testtopic/1; String topicTestzmtest; String data hello MQTT test; // mqttConfig.publish(topicTest, data); // // 订阅 mqttConfig.subscribe(topicTest); // // 发布消息 // mqttConfig.publish(pubTopic, data); // 断开连接 // mqttConfig.disconnect(); } }6、注意下面是高级版的重连和重新订阅1设置自动重连//是否自动重连 connOpts.setAutomaticReconnect(true);2不用新建回调类直接写client.setCallback(new MqttCallbackExtended() { Override public void connectComplete(boolean reconnect, String serverURI) { // 连接/重连成功后订阅 System.out.println(serverURI); try { client.subscribe(test1,qos); client.subscribe(test2,qos); } catch (Exception e) { e.printStackTrace(); } } // 连接丢失后一般在这里面进行重连 Override public void connectionLost(Throwable throwable) { System.out.println(连接丢失1); } Override public void messageArrived(String topic, MqttMessage message) { try { System.out.println(接收消息内容: new String(message.getPayload())); }catch (Exception ex){ ex.printStackTrace(); } } Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { System.out.println(111111); } });3完整MqttConfig配置类package com.example.superior_conjuncte_iot_rtspvideo.system_config.config; import com.example.superior_conjuncte_iot_rtspvideo.system_config.handle.OnMessageCallback; import jakarta.annotation.PostConstruct; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.eclipse.paho.client.mqttv3.*; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import java.util.Properties; Component Configuration Slf4j public class MqttConfig { Resource Lazy private OnMessageCallback onMessageCallback; Value(${spring.mqtt.username}) private String username; Value(${spring.mqtt.password}) private String password; Value(${spring.mqtt.url}) private String hostUrl; Value(${spring.mqtt.client-id}) private String clientId; Value(${spring.mqtt.timeout}) private Integer timeout; Value(${spring.mqtt.keep-alive}) private Integer keepAlive; Value(${spring.mqtt.qos}) private Integer qos; Value(${kafka.kafkaTopic}) private String kafkaTopic; private String topiccamera/1/video; Resource private KafkaAttributeConfig kafkaAttributeConfig; private KafkaProducerString, String kafkaVideoProducernull; private MqttClient client; /** * 项目启动时自动连接 MQTT */ PostConstruct public void init() { //初始化kafka initKafkaProducer(); connect(); } public void initKafkaProducer(){ Properties propskafkaAttributeConfig.initProductConfig(); this.kafkaVideoProducer new KafkaProducer(props); } //发送到kafka生产者 public void sendKafkaVideo(String key,String message){ ProducerRecordString,String record new ProducerRecord(kafkaTopic,key,message); kafkaVideoProducer.send(record); } //断线手动重连重新订阅 public void reConnectSubscribe() { while (true){ try { log.warn(开始重连重订阅); Thread.sleep(2000); this.connect(); // //订阅 // this.client.subscribe(myTest,2); break; }catch (Exception ex){ ex.printStackTrace(); } } } /** * 连接 MQTT */ public void connect() { try { client new MqttClient(hostUrl, clientId, new MemoryPersistence()); // MQTT 连接选项 MqttConnectOptions connOpts new MqttConnectOptions(); connOpts.setUserName(username); connOpts.setPassword(password.toCharArray()); // 保留会话 connOpts.setCleanSession(true); // 设置超时时间单位秒 connOpts.setConnectionTimeout(timeout); // 设置心跳时间单位秒表示服务器每隔1.5*20秒的时间向客户端发送心跳判断客户端是否在线 connOpts.setKeepAliveInterval(keepAlive); //是否自动重连 connOpts.setAutomaticReconnect(true); // 设置回调 // client.setCallback(onMessageCallback); client.setCallback(new MqttCallbackExtended() { Override public void connectComplete(boolean reconnect, String serverURI) { // 连接/重连成功后订阅 System.out.println(serverURI); try { client.subscribe(test1,qos); client.subscribe(test2,qos); } catch (Exception e) { e.printStackTrace(); } } // 连接丢失后一般在这里面进行重连 Override public void connectionLost(Throwable throwable) { System.out.println(连接丢失1); } Override public void messageArrived(String topic, MqttMessage message) { try { System.out.println(接收消息内容: new String(message.getPayload())); }catch (Exception ex){ ex.printStackTrace(); } } Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { System.out.println(111111); } }); // 建立连接 client.connect(connOpts); //订阅 // client.subscribe(topic,qos); // client.subscribe(test1,qos); // client.subscribe(test2,qos); } catch (MqttException me) { log.error(MQTT连接失败{}, me.getMessage()); } } /** * 断开连接 */ public void disconnect() { try { client.disconnect(); client.close(); } catch (MqttException ex) { log.error(MQTT断开连接失败{}, ex.getMessage()); } } }
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进