EMQX 是一个高度可扩展的企业级 MQTT 消息服务器,支持大规模的 IoT 设备连接和消息传输。它不仅提供了稳定的消息服务,还具有丰富的插件生态,便于集成到现有的系统中。本文将介绍如何使用 EMQX 与 Java 结合开发物联网应用。

准备工作

安装 EMQX

  1. 下载 EMQX:根据你的操作系统从 EMQX 官网 下载合适的版本。
  2. 启动 EMQX:解压后进入 bin 目录执行 ./emqx start 启动 EMQX 服务。默认情况下,EMQX 的 Dashboard 可以通过访问 http://localhost:18083 来管理,默认用户名和密码为 admin 和 public

Maven 依赖

为了在 Java 应用中使用 MQTT 协议,我们需要引入 Eclipse Paho 库:

 

xml

深色版本

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

创建 MQTT 客户端

下面是一个简单的例子,展示如何创建一个 MQTT 客户端来订阅主题并发布消息至 EMQX。

订阅主题示例

 

java

深色版本

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

public class MqttClientExample {
    public static void main(String[] args) throws MqttException {
        String broker = "tcp://localhost:1883";
        String clientId = "JavaSample";
        MemoryPersistence persistence = new MemoryPersistence();

        MqttClient client = new MqttClient(broker, clientId, persistence);
        MqttConnectOptions options = new MqttConnectOptions();
        options.setCleanSession(true);

        System.out.println("Connecting to broker: " + broker);
        client.connect(options);
        System.out.println("Connected");

        // 订阅主题
        client.subscribe("iot/device/temperature", (topic, msg) -> {
            System.out.println("\nMessage arrived. Topic: " + topic + " Message: " + new String(msg.getPayload()));
        });

        // 保持客户端运行
        Thread.sleep(30000);
        client.disconnect();
        System.out.println("Disconnected");
    }
}

发布消息示例

 

java

深色版本

import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.MqttMessage;

public class MqttPublisherExample {
    public static void main(String[] args) throws MqttException {
        String broker = "tcp://localhost:1883";
        String clientId = "Publisher";
        MemoryPersistence persistence = new MemoryPersistence();

        MqttClient client = new MqttClient(broker, clientId, persistence);
        MqttConnectOptions options = new MqttConnectOptions();
        options.setCleanSession(true);

        client.connect(options);
        System.out.println("Connected to broker: " + broker);

        String content = "Hello from Publisher!";
        MqttMessage message = new MqttMessage(content.getBytes());
        message.setQos(2); // 设置服务质量等级

        client.publish("iot/device/temperature", message);
        System.out.println("Message published");

        client.disconnect();
        System.out.println("Disconnected");
    }
}

集成 Spring Boot

为了让我们的应用更加模块化和易于维护,可以将其集成到 Spring Boot 中。

添加必要的依赖

除了之前的 Paho MQTT 客户端库外,还需要添加 Spring Boot Starter Web 依赖:

 

xml

深色版本

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>

创建 MQTT 连接配置类

 

java

深色版本

import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;

@Configuration
public class MqttConfig {

    @Bean
    public MqttPahoClientFactory mqttClientFactory() {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        MqttConnectOptions options = new MqttConnectOptions();
        options.setServerURIs(new String[]{"tcp://localhost:1883"});
        factory.setConnectionOptions(options);
        return factory;
    }
}

实现消息监听器

 

java

深色版本

import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.messaging.Message;
import org.springframework.stereotype.Component;

@Component
public class MqttMessageListener {

    private static final Logger logger = LoggerFactory.getLogger(MqttMessageListener.class);

    @ServiceActivator(inputChannel = "mqttInputChannel")
    public void handleMessage(Message<byte[]> message) {
        MqttMessage mqttMessage = (MqttMessage) message.getPayload();
        logger.info("Received MQTT message: {}", new String(mqttMessage.getPayload()));
    }
}

以上就是使用 EMQX 和 Java 开发物联网应用的基本流程。通过这种方式,你可以轻松地构建起设备与云端之间的通信桥梁,并进一步实现数据处理、存储等功能

Logo

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

更多推荐