1. mqtt 协议介绍

     了解MQTT协议以及利用Netty搭建MQTT服务器_Nanaw_z的博客-CSDN博客_mqtt netty 

     协议介绍 | EMQ X 企业版

2. mqtt broker搭建

以EMQX 社区开源版为例(企业版收费)

2.1 下载 EMQX

2.2 emqx 启动

emqx.cmd start

2.3 浏览器访问emqx 控制台面板:http://localhost:18083/#/rules

2.4 默认登录账号: admin ,密码: public

 2.5 进入面板后可能会出现:url not found 的错误,解决:

# 1. 进入emqx 目录,例如:C:\E\emqx-windows-4.3.8\emqx\etc\plugins

# 2. 修改 emqx_management.conf 文件,找到:management.listener.http 将端口更改未被占用的端口

2.6 配置主题topic

2.7 启动后监听端口:1883 

3. mqtt.fx 测试

       下载地址戳这: 官网下载地址

3.1  设置连接参数:

连接地址就是上面emqx 的地址

3.2 发布或者订阅

4. SDK 操作

4.1 引入依赖

<dependency>
    <groupId>org.eclipse.paho</groupId>
    <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
    <version>1.2.2</version>
</dependency>

4.2 测试代码

package com.cn.web.es;

import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;

/**
 * description: TestMqtt <br>
 * date: 22.2.17 14:18 <br>
 * author: cn_yaojin <br>
 * version: 1.0 <br>
 */
public class TestMqtt {

    public static void main(String[] args) {
        String t = "home/garden/fountain";
        String content = "{\"id\":\"123123\",\"accountName\":\"cn_yang\"}";
        int qos = 2;
        String broker = "tcp://127.0.0.1:1883";
        String clientId = "emqx_test";
        MemoryPersistence persistence = new MemoryPersistence();

        try {
            MqttClient client = new MqttClient(broker, clientId, persistence);

            // MQTT 连接选项
            MqttConnectOptions connOpts = new MqttConnectOptions();
//            connOpts.setUserName("emqx_test");
//            connOpts.setPassword("emqx_test_password".toCharArray());
            // 保留会话
            connOpts.setCleanSession(true);

            // 设置回调
            client.setCallback(new OnMessageCallback());

            // 建立连接
            System.out.println("Connecting to broker: " + broker);
            client.connect(connOpts);

            System.out.println("Connected");
            System.out.println("Publishing message: " + content);

            // 订阅
            client.subscribe(t);

            // 消息发布所需参数
            MqttMessage message = new MqttMessage(content.getBytes());
//            message.setQos(qos);
//            client.publish(t, message);
//            System.out.println("Message published");

//            client.disconnect();
//            client.close();
        } 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();
        }

    }

    public static class OnMessageCallback implements MqttCallback {

        @Override
        public void connectionLost(Throwable cause) {
            // 连接丢失后,一般在这里面进行重连
            System.out.println("连接断开,可以做重连");
        }

        @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());
        }
    }

}

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐