一:在pom文件中添加依赖

<!--添加MQTT依赖-->
        <dependency>
            <groupId>org.eclipse.paho</groupId>
            <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
            <version>1.2.5</version>
        </dependency>

二:新建一个MqttConfig类

import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttCallback;
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.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.stereotype.Component;

@Slf4j
@Configuration
public class MqttConfig {

    private static final String MQTT_BROKER_URL = "tcp://127.0.0.1:1883";  // 替换为自己的 MQTT 服务器地址
    private static final String CLIENT_ID = "yu_mqtt_client";//客户端ID,一般唯一
    private static final String USERNAME = "ytong";  // 替换为自己的用户名
    private static final String PASSWORD = "admin123";  // 替换为自己的密码


    @Bean
    public MqttClient mqttClient() throws MqttException {
        MqttClient client = new MqttClient(MQTT_BROKER_URL, CLIENT_ID, new MemoryPersistence());
        MqttConnectOptions options = new MqttConnectOptions();

        // 设置用户名和密码
        options.setUserName(USERNAME);
        options.setPassword(PASSWORD.toCharArray());

        options.setAutomaticReconnect(true); // 设置自动重连
        options.setCleanSession(false); // 设置会话不被清除,保持连接状态
        options.setConnectionTimeout(10); // 设置连接超时时间
        options.setKeepAliveInterval(60); // 设置会话心跳时间

        // 设置回调
        client.setCallback(new MqttCallback() {
            @Override
            public void connectionLost(Throwable cause) {
                log.info("连接丢失!原因:" + cause.getMessage());
            }
            @Override
            public void messageArrived(String topic, MqttMessage message) throws Exception {
                log.info("接收到消息:\n主题:" + topic + "\n消息:" + new String(message.getPayload()));
            }

            @Override
            public void deliveryComplete(org.eclipse.paho.client.mqttv3.IMqttDeliveryToken token) {
                int messageId = token.getMessageId();
                String[] topics = token.getTopics();
                log.info("消息发送成功!id为{},主题为{}",messageId,topics);
            }
        });
        client.connect(options);
        log.info("MQTT服务器连接成功!");
        return client;
    }
}

三:新建一个MqttService类(里面定义两个方法,一个订阅,一个发布)

import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

@Service
public class MqttService {

    @Autowired
    private MqttClient mqttClient;

    //向某一个主题发布消息
    public void publish(String topic, String message) throws MqttException {
        MqttMessage mqttMessage = new MqttMessage(message.getBytes());
        mqttMessage.setQos(1);
        mqttClient.publish(topic, mqttMessage);
    }

    //订阅某一个主题
    public void subscribe(String topic) throws MqttException {
        mqttClient.subscribe(topic);
    }
}

三:新建一个MQTT Controller接口,用来测试

import com.fasterxml.jackson.core.JsonProcessingException;
import com.wanhe.common.core.controller.BaseController;
import com.wanhe.power.service.MqttService;
import com.wanhe.power.utils.MqttSendMsgUtils;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping("/system/mqtt")
public class MqttController extends BaseController {

    @Autowired
    private MqttService mqttService;

    @PostMapping("/send")
    public String sendMessage() {
        String topic = "device/status";//要发送的主题
        
        try {
            String message = "这里填写你发送的消息!";
//          String message  = MqttSendMsgUtils.sendMqttMsg(1);
            mqttService.publish(topic, message);
        } catch (Exception e) {
            logger.error("发送格式组合失败,错误原因为{}",e.getMessage());
        }
        return "主题发送到: " + topic;
    }

    @PostMapping("/subscribe")
    public String subscribeToTopic(@RequestParam String topic) {
        try {
            mqttService.subscribe(topic);
            return "Subscribed to topic: " + topic;
        } catch (MqttException e) {
            e.printStackTrace();
            return "Error subscribing to topic";
        }
    }
}

四:工具推荐(本地测试)

1.EMQX服务器
下载EMQX

EMQX官方网站 官网下载地址提供了EMQX开源版软件包
但是5.3.2版本以后就没有提供Windows系统软件包,需要到这里下载Windows下载地址,选择5.3.2之前的版本
2.MQTTX 客户端
推荐:如果不会,请参照这篇博客,里面很详细

总结:

这种只适用于简单的应用场景,对于复杂的物联网场景,这种方式可能不是最好的选择。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐