mqtt 服务搭建及测试
·
1. mqtt 协议介绍
了解MQTT协议以及利用Netty搭建MQTT服务器_Nanaw_z的博客-CSDN博客_mqtt netty
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());
}
}
}
更多推荐
所有评论(0)