下面为你提供一个连接到 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 进行了优化:

  1. 增加了完整的回调函数实现,能更好地处理 EMQX 的连接状态和消息确认
  2. 设置了适合 EMQX 的心跳间隔(60秒)
  3. 使用了带时间戳的客户端ID,避免在EMQX中出现客户端ID冲突
  4. 支持EMQX的保留消息功能(通过message.setRetained()设置)
  5. 完善了错误处理机制,方便排查连接问题

使用前请注意:

  • 确保 EMQX 5.3.2 服务器已正确安装并启动
  • 在 EMQX 控制台中已创建用户名为 “user”、密码为 “123456” 的账号
  • 根据你的 EMQX 服务器地址修改 BROKER 常量(默认是本地服务器)
  • 如果 EMQX 配置了 SSL/TLS,需要将连接协议改为 “ssl://” 并使用相应端口(通常是8883)

运行程序后,如果看到"连接EMQX成功"和"消息发送成功"的提示,说明已成功与 EMQX 服务器通信。你可以在 EMQX 控制台的"监控"或"消息"部分查看发送的测试消息。

更多推荐