一、什么是MQTT

        MQTT(Message Queuing Telemetry Transport)作为一种轻量级的物联网通信协议,在设备间低带宽、高延迟的场景下表现出色。主要由发布者、服务器、订阅者共同来实现。

发布者(Publish):可以理解为消息的发送方

订阅者(Subscribe):可以理解为消息的接收方

服务器(Broker):可以理解为一个中转站,或者一个容器

消息结构组成:

  1. 主题(Topic):消息的分类标识,采用层级结构。
  2. 负载(Payload):消息的内容,通常是String类型。

        当发布者发布一个消息到到服务器,服务器发现该消息的主题与订阅者订阅的主题相同则会发送这个消息给订阅者,它不用想Http协议一样,如果想实现一个点对点的通信必须知道ip,它可以通过服务器来将消息发送给订阅了相同主题的订阅者。

二、Java 客户端连接 EMQX

1、添加依赖

<dependencies>
    <dependency>
        <groupId>org.eclipse.paho</groupId>
        <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
        <version>1.2.5</version>
    </dependency>
    <!-- 用于JSON处理 -->
    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>fastjson</artifactId>
        <version>1.2.83</version>
    </dependency>
</dependencies>

2、极简 “发送者”(对接 EMQX,上报电量)

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

public class Sender { // 发送者:柜机上报电量
    public static void main(String[] args) throws MqttException {
        // 1. 连接 EMQX 服务器(Broker 地址:tcp://localhost:1883)
        MqttClient client = new MqttClient("tcp://localhost:1883", "sender-cabinet-001", new MemoryPersistence());
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName("admin"); // EMQX 默认用户名
        options.setPassword("public".toCharArray()); // EMQX 默认密码
        client.connect(options);

        // 2. 发送消息(主题:cabinet/001/battery,内容:80%)
        String topic = "cabinet/001/battery";
        String content = "80%";
        MqttMessage message = new MqttMessage(content.getBytes());
        message.setQos(1); // 服务质量
        client.publish(topic, message);

        client.disconnect();
    }
}

3. 极简 “订阅者”(对接 EMQX,接收电量)

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

public class Subscriber { // 订阅者:后端系统接收电量
    public static void main(String[] args) throws MqttException {
        // 1. 连接 EMQX 服务器
        MqttClient client = new MqttClient("tcp://localhost:1883", "subscriber-backend-001", new MemoryPersistence());
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName("admin");
        options.setPassword("public".toCharArray());
        
        // 2. 接收消息的回调(收到 EMQX 转发的消息后执行)
        client.setCallback(new MqttCallback() {
            @Override
            public void messageArrived(String topic, MqttMessage message) throws Exception {
                System.out.println("收到消息:主题=" + topic + ",内容=" + new String(message.getPayload()));
            }
            @Override public void connectionLost(Throwable cause) {}
            @Override public void deliveryComplete(IMqttDeliveryToken token) {}
        });

        // 3. 订阅主题(监听 cabinet/001 的电量)
        client.connect();
        client.subscribe("cabinet/001/battery", 1);
    }
}

三、总结

        当然你也可以对EMQX服务器进行一个“EMQX专用客户端”的封装。它可 “让你更方便地实现发送者 / 订阅者”,减少代码的冗余。在实际业务中我们可以还会引用MQ这个消息中间件,能解决设备通信、高并发、业务解耦等核心问题。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐