【MQTT】【实战】---- java编写单元测试给mqtt服务器发送消息
·
下面为你提供一个连接到 EMQX 5.3.2 服务器并发送测试消息的 Spring Boot 程序。
EMQX 是一个开源的 MQTT 消息服务器,代码会针对 EMQX 的特性进行适当配置。
首先确保 pom.xml 中有所需依赖(与之前相同,包含 Spring Boot 和 Paho MQTT 客户端):
<dependencies>
<!-- Spring Boot Starter -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<!-- Eclipse Paho MQTT Client -->
<dependency>
<groupId>org.eclipse.paho</groupId>
<artifactId>org.eclipse.paho.client.mqttv3</artifactId>
<version>1.2.5</version>
</dependency>
</dependencies>
然后是主程序代码:
package com.zgb.app.test.mqtt;
import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class EmqxTestClient {
// EMQX服务器地址,默认端口1883(TCP),如果启用了SSL则使用8883
private static final String BROKER = "tcp://192.168.0.71:1883";
// 客户端ID,EMQX建议使用有意义的唯一标识
private static final String CLIENT_ID = "spring-boot-emqx-test-" + System.currentTimeMillis();
// 要发送消息的主题,EMQX支持多级主题
private static final String TOPIC = "emqx/test/topic";
// QoS级别,EMQX完全支持QoS 0, 1, 2
private static final int QOS = 1;
// EMQX中配置的用户名和密码
private static final String USERNAME = "deviceUser";
private static final String PASSWORD = "szty123";
public static void main(String[] args) {
// 内存持久化
MemoryPersistence persistence = new MemoryPersistence();
try {
// 创建MQTT客户端
MqttClient client = new MqttClient(BROKER, CLIENT_ID, persistence);
// 配置连接选项,针对EMQX进行优化
MqttConnectOptions connOpts = new MqttConnectOptions();
connOpts.setUserName(USERNAME);
connOpts.setPassword(PASSWORD.toCharArray());
// EMQX支持保持连接心跳,设置为60秒
connOpts.setKeepAliveInterval(60);
// 自动重连
connOpts.setAutomaticReconnect(true);
// 连接超时
connOpts.setConnectionTimeout(10);
// 清除会话设置,根据需求调整
connOpts.setCleanSession(true);
// 设置回调函数,处理连接和消息相关事件
client.setCallback(new MqttCallback() {
@Override
public void connectionLost(Throwable cause) {
System.out.println("连接丢失: " + cause.getMessage());
}
@Override
public void messageArrived(String topic, MqttMessage message) throws Exception {
System.out.println("收到消息 - 主题: " + topic + ", 内容: " + new String(message.getPayload()));
}
@Override
public void deliveryComplete(IMqttDeliveryToken token) {
System.out.println("消息已送达,消息ID: " + token.getMessageId());
}
});
System.out.println("连接到EMQX服务器: " + BROKER);
// 连接到EMQX服务器
client.connect(connOpts);
System.out.println("连接EMQX成功,客户端ID: " + CLIENT_ID);
// 要发送的测试消息
String content = "这是一条发送到EMQX 5.3.2的测试消息: " + System.currentTimeMillis();
System.out.println("准备发送消息: " + content);
// 创建消息对象
MqttMessage message = new MqttMessage(content.getBytes());
message.setQos(QOS);
// 保留消息设置,EMQX支持保留消息功能
message.setRetained(false);
// 发布消息到指定主题
client.publish(TOPIC, message);
System.out.println("消息发送成功,主题: " + TOPIC + ", QoS: " + QOS);
// 等待消息确认
Thread.sleep(1000);
// 断开连接
client.disconnect();
System.out.println("已与EMQX服务器断开连接");
// 关闭客户端
client.close();
System.exit(0);
} catch (MqttException me) {
System.out.println("MQTT错误: " + me.getMessage());
System.out.println("错误代码: " + me.getReasonCode());
System.out.println("错误原因: " + me.getCause());
me.printStackTrace();
} catch (InterruptedException e) {
System.out.println("线程中断: " + e.getMessage());
Thread.currentThread().interrupt();
}
}
}
测试:

这个程序针对 EMQX 5.3.2 进行了优化:
- 增加了完整的回调函数实现,能更好地处理 EMQX 的连接状态和消息确认
- 设置了适合 EMQX 的心跳间隔(60秒)
- 使用了带时间戳的客户端ID,避免在EMQX中出现客户端ID冲突
- 支持EMQX的保留消息功能(通过message.setRetained()设置)
- 完善了错误处理机制,方便排查连接问题
使用前请注意:
- 确保 EMQX 5.3.2 服务器已正确安装并启动
- 在 EMQX 控制台中已创建用户名为 “user”、密码为 “123456” 的账号
- 根据你的 EMQX 服务器地址修改 BROKER 常量(默认是本地服务器)
- 如果 EMQX 配置了 SSL/TLS,需要将连接协议改为 “ssl://” 并使用相应端口(通常是8883)
运行程序后,如果看到"连接EMQX成功"和"消息发送成功"的提示,说明已成功与 EMQX 服务器通信。你可以在 EMQX 控制台的"监控"或"消息"部分查看发送的测试消息。
更多推荐



所有评论(0)