Programación de Clientes MQTT en Vue.js y Java

Este artículo detalla cómo implementar clientes MQTT utilizando Vue.js en el frontend y Java (con Eclipse Paho y Spring Integration) en el backend, junto con un caso de uso práctico para controlar bombillas inteligentes.

  1. Implementación de Cliente MQTT en Vue.js

1.1. Configuración Inicial del Proyecto

Comience creando un nuevo proyecto Vue.js usando Vite e instale las dependencias necesarias, incluyendo Element Plus para la interfaz de usuario y mqtt.js para la comunicación MQTT.


# Crear proyecto Vue.js con Vite
npm create vite@latest my-mqtt-app --template vue
cd my-mqtt-app

# Instalar dependencias del proyecto
npm install

# Instalar Element Plus y mqtt.js
npm install element-plus mqtt --save

Modifique main.js para usar Element Plus:


import { createApp } from 'vue'
import App from './App.vue'
import ElementPlus from 'element-plus'
import 'element-plus/dist/index.css'

const app = createApp(App)
app.use(ElementPlus)
app.mount('#app')

1.2. Creación del Componente MqttDemo.vue

Cree un nuevo componente MqttDemo.vue en la carpeta components para gestionar la interfaz y la lógica de MQTT.

<script setup>
import { ref } from "vue";
import mqtt from "mqtt";

// Opciones de Calidad de Servicio (QoS)
const qosList = [0, 1, 2];

// Estado de la conexión
const connectionInfo = ref({
  protocol: 'ws',
  host: "192.168.136.147",
  port: 8083,
  clientId: "emqx_vue3_" + Math.random().toString(16).substring(2, 8),
  username: "zhangsan",
  password: "123",
  clean: true,
  connectTimeout: 10 * 1000, // ms
  reconnectPeriod: 4000,     // ms
});

const client = ref({});
const isConnected = ref(false);

// Estado de la suscripción
const subscriptionInfo = ref({
  topic: '',
  qos: 0
});
const isSubscribed = ref(false);

// Estado de la publicación
const publishInfo = ref({
  topic: '',
  qos: 0,
  payload: ''
});

// Mensajes recibidos
const receivedMessages = ref("");

// Funciones de conexión
const createConnection = () => {
  const { protocol, host, port, ...options } = connectionInfo.value;
  const connectUrl = `${protocol}://${host}:${port}/mqtt`;
  console.log(`Conectando a: ${connectUrl}`);
  client.value = mqtt.connect(connectUrl, options);
  isConnected.value = true;
  console.info("Conexión establecida...");

  client.value.on('message', (topic, message) => {
    console.info(`Mensaje recibido: Tópico=${topic}, Mensaje=${message}`);
    receivedMessages.value = `Tópico: ${topic}, Mensaje: ${message}`;
  });

  client.value.on('error', (err) => {
    console.error("Error de conexión:", err);
    isConnected.value = false;
  });

  client.value.on('close', () => {
    console.info("Conexión cerrada.");
    isConnected.value = false;
  });
};

const closeConnection = () => {
  if (client.value && client.value.end) {
    client.value.end(false, () => {
      isConnected.value = false;
      console.info("Conexión cerrada exitosamente...");
    });
  }
};

// Funciones de suscripción
const subscribeTopicHandler = () => {
  const { topic, qos } = subscriptionInfo.value;
  if (client.value && client.value.subscribe) {
    client.value.subscribe(topic, { qos }, (error, res) => {
      if (error) {
        console.info("Error al suscribir:", error);
        return;
      }
      isSubscribed.value = true;
      console.info(`Suscrito a ${topic} con QoS ${qos}`);
    });
  }
};

const unsubscribeTopicHandler = () => {
  const { topic } = subscriptionInfo.value;
  if (client.value && client.value.unsubscribe) {
    client.value.unsubscribe(topic, (error, res) => {
      if (error) {
        console.info("Error al cancelar suscripción:", error);
        return;
      }
      isSubscribed.value = false;
      console.info(`Cancelada suscripción a ${topic}`);
    });
  }
};

// Funciones de publicación
const doPublish = () => {
  const { topic, payload, qos } = publishInfo.value;
  if (client.value && client.value.publish) {
    client.value.publish(topic, payload, { qos }, (error) => {
      if (error) {
        console.info("Error al publicar mensaje:", error);
        return;
      }
      console.info(`Mensaje publicado: Tópico=${topic}, Payload=${payload}, QoS=${qos}`);
    });
  }
};
</script>

<template>
  <div class="mqtt-demo">
    <el-card>
      <h1>Información de Conexión</h1>
      <el-form label-position="top" >
        <el-row :gutter="20">
          <el-col :span="8">
            <el-form-item label="Protocolo">
              <el-select v-model="connectionInfo.protocol">
                <el-option label="ws://" value="ws"></el-option>
                <el-option label="wss://" value="wss"></el-option>
              </el-select>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="Host">
              <el-input v-model="connectionInfo.host"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="Puerto">
              <el-input type="number" v-model="connectionInfo.port" placeholder="8083/8084"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="ClientID">
              <el-input v-model="connectionInfo.clientId"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="Usuario">
              <el-input v-model="connectionInfo.username"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="Contraseña">
              <el-input v-model="connectionInfo.password" type="password"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="24">
            <el-button type="primary" :disabled="isConnected" @click="createConnection">Conectar</el-button>
            <el-button type="danger" :disabled="!isConnected" @click="closeConnection">Desconectar</el-button>
          </el-col>
        </el-row>
      </el-form>
    </el-card>

    <el-card>
      <h1>Suscripción a Tópico</h1>
      <el-form label-position="top" >
        <el-row :gutter="20">
          <el-col :span="8">
            <el-form-item label="Tópico">
              <el-input v-model="subscriptionInfo.topic"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="QoS">
              <el-select v-model="subscriptionInfo.qos">
                <el-option
                    v-for="qos in qosList"
                    :key="qos"
                    :label="qos"
                    :value="qos"
                ></el-option>
              </el-select>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-button type="primary" class="sub-btn"
                       :disabled="isSubscribed || !isConnected"
                       @click="subscribeTopicHandler">Suscribir</el-button>
            <el-button type="primary" class="sub-btn"
                       :disabled="!isSubscribed"
                       @click="unsubscribeTopicHandler">Cancelar Suscripción</el-button>
          </el-col>
        </el-row>
      </el-form>
    </el-card>

    <el-card>
      <h1>Publicar Mensaje</h1>
      <el-form label-position="top" >
        <el-row :gutter="20">
          <el-col :span="8">
            <el-form-item label="Tópico">
              <el-input v-model="publishInfo.topic"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="Payload">
              <el-input v-model="publishInfo.payload"></el-input>
            </el-form-item>
          </el-col>
          <el-col :span="8">
            <el-form-item label="QoS">
              <el-select v-model="publishInfo.qos">
                <el-option
                    v-for="qos in qosList"
                    :key="qos"
                    :label="qos"
                    :value="qos"
                ></el-option>
              </el-select>
            </el-form-item>
          </el-col>
        </el-row>
      </el-form>
      <el-col :span="24" class="text-right">
        <el-button type="primary" :disabled="!isConnected" @click="doPublish">Publicar</el-button>
      </el-col>
    </el-card>

    <el-card>
      <h1>Mensajes Recibidos</h1>
      <el-col :span="24">
        <el-input
            type="textarea"
            :rows="3"
            readonly
            v-model="receivedMessages"
        ></el-input>
      </el-col>
    </el-card>
  </div>
</template>

<style>
.mqtt-demo {
  max-width: 1200px;
  margin: 32px auto 0 auto;
}
h1 {
  font-size: 16px;
  margin-top: 0;
}
.el-card {
  margin-bottom: 32px;
}
.el-card__body {
  padding: 24px;
}
.el-select {
  width: 100%;
}
.text-right {
  text-align: right;
}
.sub-btn {
  margin-top: 30px;
}
</style>

1.3. Actualización de App.vue

Modifique App.vue para incluir el componente MqttDemo.

<script setup>
import MqttDemo from "./components/MqttDemo.vue";
</script>

<template>
  <MqttDemo/>
</template>

<style>
/* Estilos globales si son necesarios */
</style>

  1. Implementación de Cliente MQTT en Java

2.1. Usando Eclipse Paho Java Client

Este enfoque utiliza la librería oficial org.eclipse.paho.client.mqttv3.

2.1.1. Configuración del Proyecto Spring Boot

Agregue la siguiente dependencia a su archivo pom.xml:

<dependencies>
    <!-- Dependencias de Spring Boot Test -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-test</artifactId>
        <scope>test</scope>
    </dependency>

    <!-- Cliente MQTT Eclipse Paho -->
    <dependency>
        <groupId>org.eclipse.paho</groupId>
        <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
        <version>1.2.5</version>
    </dependency>
</dependencies>

2.1.2. Conexión al Broker MQTT

Implemente la lógica para establecer una conexión.

import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.junit.jupiter.api.Test;

public class MqttConnectionTest {

    @Test
    public void createConnection() throws MqttException {
        String brokerUrl = "tcp://192.168.136.147:1883";
        String clientId = "mqtt_java_client_01";
        String username = "zhangsan";
        String password = "123";

        // Cliente MQTT con persistencia en memoria
        MqttClient client = new MqttClient(brokerUrl, clientId, new MemoryPersistence());

        // Opciones de conexión
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName(username);
        options.setPassword(password.toCharArray());
        options.setCleanSession(true); // Iniciar una sesión limpia
        options.setConnectionTimeout(60); // Tiempo de espera de conexión en segundos
        options.setKeepAliveInterval(60); // Intervalo de KeepAlive en segundos

        // Conectar al broker
        client.connect(options);
        System.out.println("Conectado al broker: " + brokerUrl);

        // Mantener la conexión activa (ejemplo simple, en una app real se maneja mejor)
        // while (!Thread.currentThread().isInterrupted()) {
        //     // Esperar o realizar otras operaciones
        // }

        // Desconexión y cierre (idealmente al finalizar la aplicación)
        // client.disconnect();
        // client.close();
    }
}

2.1.3. Publicación de Mensajes

Código para anviar un mensaje a un tópico específico.

import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.junit.jupiter.api.Test;

public class MqttPublishTest {

    @Test
    public void sendMessage() throws MqttException, InterruptedException {
        String brokerUrl = "tcp://192.168.136.147:1883";
        String clientId = "mqtt_java_publisher";
        String username = "zhangsan";
        String password = "123";
        String topic = "a/c";
        String content = "¡Hola desde Java MQTT!";

        MqttClient client = new MqttClient(brokerUrl, clientId, new MemoryPersistence());
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName(username);
        options.setPassword(password.toCharArray());
        client.connect(options);
        System.out.println("Conectado para publicar.");

        // Crear el mensaje MQTT
        MqttMessage message = new MqttMessage(content.getBytes());
        message.setQos(2); // Calidad de servicio 2
        message.setRetained(false); // No retener el mensaje en el broker

        // Publicar el mensaje
        IMqttDeliveryToken token = client.publish(topic, message);
        token.waitForCompletion(); // Esperar a que la entrega se complete

        System.out.println("Mensaje publicado en tópico '" + topic + "'");

        client.disconnect();
        client.close();
        System.out.println("Desconectado.");
    }
}

2.1.4. Suscripción y Recepción de Mensajes

Código para suscribirse a un tópico y manejar mensajes entrantes.

import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.junit.jupiter.api.Test;

public class MqttSubscribeTest {

    @Test
    public void receiveMessage() throws MqttException, InterruptedException {
        String brokerUrl = "tcp://192.168.136.147:1883";
        String clientId = "mqtt_java_subscriber_02";
        String username = "zhangsan";
        String password = "123";
        String topicToSubscribe = "a/d";
        int qosLevel = 2;

        MqttClient client = new MqttClient(brokerUrl, clientId, new MemoryPersistence());
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName(username);
        options.setPassword(password.toCharArray());
        options.setCleanSession(true);

        // Configurar el callback para manejar eventos
        client.setCallback(new MqttCallback() {
            @Override
            public void connectionLost(Throwable cause) {
                System.out.println("Conexión perdida: " + cause.getMessage());
                // Aquí se podría intentar reconectar
            }

            @Override
            public void messageArrived(String topic, MqttMessage message) {
                System.out.println("------------------------------------");
                System.out.println("Tópico recibido: " + topic);
                System.out.println("QoS: " + message.getQos());
                System.out.println("Mensaje: " + new String(message.getPayload()));
                System.out.println("------------------------------------");
            }

            @Override
            public void deliveryComplete(IMqttDeliveryToken token) {
                System.out.println("Entrega completa para el token: " + token.getMessageId());
            }
        });

        // Conectar y suscribirse
        client.connect(options);
        System.out.println("Conectado y suscrito al tópico: " + topicToSubscribe + " con QoS " + qosLevel);
        client.subscribe(topicToSubscribe, qosLevel);

        // Mantener la aplicación corriendo para recibir mensajes
        while (!Thread.currentThread().isInterrupted()) {
            Thread.sleep(1000);
        }

        // Desconexión y cierre
        client.disconnect();
        client.close();
        System.out.println("Desconectado.");
    }
}

2.2. Usando Spring Integration MQTT

Este enfoque se integra con Spring Boot para simplificar la configuración y el manejo de mensajes.

2.2.1. Configuración del Proyecto y Dependencias

Agregue las siguientes dependencias a su pom.xml:

<dependencies>
    <!-- Spring Boot Web Starter -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>

    <!-- Spring Integration Core -->
    <dependency>
        <groupId>org.springframework.integration</groupId>
        <artifactId>spring-integration-core</artifactId>
    </dependency>

    <!-- Spring Integration MQTT Adapter -->
    <dependency>
        <groupId>org.springframework.integration</groupId>
        <artifactId>spring-integration-mqtt</artifactId>
        <version>5.5.13</version> <!-- Ajustar versión según sea necesario -->
    </dependency>

    <!-- Lombok para boilerplate code -->
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>

    <!-- Fastjson para serialización/deserialización JSON -->
    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>fastjson</artifactId>
        <version>1.2.83</version>
    </dependency>
</dependencies>

2.2.2. Configuración de la Aplicación

Defina las propiedades de MQTT en application.yml:

spring:
  mqtt:
    username: zhangsan
    password: 123
    url: tcp://192.168.136.147:1883
    # Client ID para suscripción
    subClientId: spring_integration_subscriber_123
    subTopic: iot/lamp/line,iot/lamp/device/status
    # Client ID para publicación
    pubClientId: spring_integration_publisher_abc

Cree una clase de propiedades para leer la configuración:

import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;

@Data
@Component
@ConfigurationProperties(prefix = "spring.mqtt")
public class MqttProperties {
    private String username;
    private String password;
    private String url;
    private String subClientId;
    private String subTopic;
    private String pubClientId;
}

Configure la fábrica de clientes MQTT:

import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.beans.factory.annotation.Autowired;
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 {

    @Autowired
    private MqttProperties mqttProperties;

    @Bean
    public MqttPahoClientFactory mqttClientFactory() {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        MqttConnectOptions options = new MqttConnectOptions();
        options.setServerURIs(new String[]{mqttProperties.getUrl()});
        options.setUserName(mqttProperties.getUsername());
        options.setPassword(mqttProperties.getPassword().toCharArray());
        options.setCleanSession(true); // Limpiar sesión al conectar

        factory.setConnectionOptions(options);
        return factory;
    }
}

2.2.3. Configuración de Entrada (Suscripción)

Defina un canal y un adaptador para recibir mensajes:

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.IntegrationComponentScan;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;

// Componente para procesar mensajes entrantes
@Component
class ReceiverMessageHandler implements MessageHandler {

    @Autowired
    private MqttProperties mqttProperties; // Necesario para comparar tópicos

    @Override
    public void handleMessage(Message> message) throws MessagingException {
        String receivedTopic = (String) message.getHeaders().get("mqtt_receivedTopic");
        Object payload = message.getPayload();
        System.out.println("Mensaje recibido - Tópico: " + receivedTopic + ", Payload: " + payload);

        // Lógica específica basada en el tópico
        if ("iot/lamp/line".equals(receivedTopic)) {
            System.out.println("Procesando mensaje de dispositivo online...");
            // Aquí iría la lógica para actualizar el estado del dispositivo
        } else if ("iot/lamp/device/status".equals(receivedTopic)) {
            System.out.println("Procesando mensaje de cambio de estado del dispositivo...");
            // Aquí iría la lógica para registrar el cambio de estado
        }
    }
}

@Configuration
@IntegrationComponentScan // Para que Spring detecte los componentes de integración
public class MqttInboundConfig {

    @Autowired
    private MqttProperties mqttProperties;
    @Autowired
    private MqttPahoClientFactory mqttClientFactory;
    @Autowired
    private ReceiverMessageHandler receiverMessageHandler; // Inyectar el handler

    @Bean
    public MessageChannel mqttInputChannel() {
        return new DirectChannel();
    }

    @Bean
    public MqttPahoMessageDrivenChannelAdapter mqttInboundAdapter() {
        String[] topics = mqttProperties.getSubTopic().split(",");
        MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter(
                mqttProperties.getUrl(),
                mqttProperties.getSubClientId(),
                mqttClientFactory,
                topics);

        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(1); // QoS para la suscripción
        adapter.setOutputChannel(mqttInputChannel()); // Enviar mensajes a este canal
        return adapter;
    }

    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel")
    public MessageHandler handler() {
        return receiverMessageHandler; // Usar el handler inyectado
    }
}

2.2.4. Configuración de Salida (Publicación)

Configure un canal y un manejador para enviar mensajes:

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.IntegrationComponentScan;
import org.springframework.integration.annotation.MessagingGateway;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;

// Gateway para enviar mensajes MQTT
@MessagingGateway(defaultRequestChannel = "mqttOutputChannel")
public interface MqttGateway {
    @Header(value = org.springframework.integration.mqtt.support.MqttHeaders.TOPIC)
    String TOPIC = "mqtt_receivedTopic"; // Alias para el header

    @Header(value = org.springframework.integration.mqtt.support.MqttHeaders.QOS)
    String QOS = "mqtt_qos";

    void sendToMqtt(String payload); // Envía al tópico por defecto

    void sendToMqtt(@Header(TOPIC) String topic, String payload); // Envía a un tópico específico

    void sendToMqtt(@Header(TOPIC) String topic, @Header(QOS) int qos, String payload); // Envía con QoS
}

// Bean para enviar mensajes
@Component
class MqttMessageSender {

    private final MqttGateway mqttGateway;

    public MqttMessageSender(MqttGateway mqttGateway) {
        this.mqttGateway = mqttGateway;
    }

    public void sendToTopic(String topic, String payload) {
        mqttGateway.sendToMqtt(topic, payload);
        System.out.println("Enviado a tópico: " + topic + ", Payload: " + payload);
    }

    public void sendToTopic(String topic, int qos, String payload) {
        mqttGateway.sendToMqtt(topic, qos, payload);
        System.out.println("Enviado a tópico: " + topic + ", QoS: " + qos + ", Payload: " + payload);
    }
}


@Configuration
public class MqttOutboundConfig {

    @Autowired
    private MqttProperties mqttProperties;
    @Autowired
    private MqttPahoClientFactory mqttClientFactory;

    @Bean
    public MessageChannel mqttOutputChannel() {
        return new DirectChannel();
    }

    @Bean
    @ServiceActivator(inputChannel = "mqttOutputChannel")
    public MessageHandler mqttOutboundMessageHandler() {
        // URL y Client ID para el publicador
        MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler(
                mqttProperties.getUrl(),
                mqttProperties.getPubClientId(),
                mqttClientFactory);

        messageHandler.setAsync(true); // Procesamiento asíncrono
        messageHandler.setDefaultQos(1); // QoS por defecto para publicación
        messageHandler.setDefaultTopic("default/topic"); // Tópico por defecto si no se especifica
        return messageHandler;
    }
}

Para usar el servicio de envío, inyecte MqttMessageSender en sus controladores o servicios.

  1. Caso de Uso: Control de Bombillas Inteligentes

Este escenario demuestra cómo usar MQTT para interactuar con dispositivos IoT.

3.1. Preparación del Entorno Backend

  1. Base de Datos: Cree las tablas tb_lamp (para estado online/offline del dispositivo) y tb_lamp_status (para estado on/off de la bombilla).
  2. Dependencias Adicionales: Añada mybatis-plus-boot-starter y el driver de MySQL a su pom.xml.
  3. Configuración de DataSource: Configure la conexión a la base de datos en application.yml.
  4. Generación de Código MyBatis-Plus: Use el generador de código para crear los Mappers y Entidades.
  5. Escaneo de Mappers: Añada la anotación @MapperScan a su clase principal de Spring Boot.

3.2. Lógica del Dispositivo (Online/Offline)

Tópico: iot/lamp/line
Mensaje (JSON): {"deviceId": "XXX", "online": 1} (1 para online, 0 para offline)

Modifique ReceiverMessageHandler para procesar este mensaje:

import com.alibaba.fastjson.JSON;
// ... otros imports
import com.example.demo.service.TbLampService; // Asumiendo que tienes esta interfaz
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.stereotype.Component;
import java.util.Map;

@Component
public class ReceiverMessageHandler implements MessageHandler {

    @Autowired
    private TbLampService tbLampService; // Servicio para gestionar lámparas
    @Autowired
    private TbLampStatusService tbLampStatusService; // Servicio para estado de bombilla

    @Override
    public void handleMessage(Message> message) throws MessagingException {
        String receivedTopic = (String) message.getHeaders().get("mqtt_receivedTopic");
        String payload = message.getPayload().toString();
        System.out.println("Mensaje Recibido - Tópico: " + receivedTopic + ", Payload: " + payload);

        if ("iot/lamp/line".equals(receivedTopic)) {
            tbLampService.updateLampOnlineStatus(payload);
        } else if ("iot/lamp/device/status".equals(receivedTopic)) {
            tbLampStatusService.saveDeviceStatus(payload);
        }
    }
}

Implemente la lógica en TbLampServiceImpl (asumiendo que usa MyBatis-Plus):

import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.example.demo.entity.TbLamp; // Entidad de lámpara
import com.example.demo.mapper.TbLampMapper; // Mapper de lámpara
import com.example.demo.service.TbLampService;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.Map;
import com.alibaba.fastjson.JSON;

@Service
public class TbLampServiceImpl extends ServiceImpl<tblampmapper tblamp=""> implements TbLampService {

    @Override
    public void updateLampOnlineStatus(String jsonInfo) {
        // Parsear el JSON para obtener deviceId y status
        Map<string object=""> map = JSON.parseObject(jsonInfo, Map.class);
        String deviceId = map.get("deviceId").toString();
        Integer status = Integer.parseInt(map.get("online").toString());

        // Buscar la lámpara por deviceId
        LambdaQueryWrapper<tblamp> queryWrapper = new LambdaQueryWrapper<>();
        queryWrapper.eq(TbLamp::getDeviceid, deviceId);
        TbLamp lamp = getOne(queryWrapper);

        if (lamp == null) {
            // Si no existe, crear una nueva entrada
            lamp = new TbLamp();
            lamp.setDeviceid(deviceId);
            lamp.setStatus(status);
            lamp.setCreateTime(new Date());
            lamp.setUpdateTime(new Date());
            save(lamp);
            System.out.println("Dispositivo nuevo registrado: " + deviceId);
        } else {
            // Si existe, actualizar el estado y la fecha de actualización
            lamp.setStatus(status);
            lamp.setUpdateTime(new Date());
            updateById(lamp);
            System.out.println("Dispositivo actualizado: " + deviceId + ", Estado: " + status);
        }
    }
}
</tblamp></string></tblampmapper>

3.3. Lógica de Control (Encender/Apagar)

Tópico: iot/lamp/server/status
Mensaje (JSON): {"deviceId": "XXX", "status": 0} (0 para apagar, 1 para encender)

Cree un endpoint REST para enviar comandos:

import com.alibaba.fastjson.JSON;
import com.example.mqtt.sender.MqttMessageSender; // Suponiendo que tienes esta clase
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;

import java.util.HashMap;
import java.util.Map;

@RestController
@RequestMapping("/api/lamp")
public class LampController {

    @Autowired
    private MqttMessageSender mqttMessageSender;

    @GetMapping("/{deviceId}/{status}")
    public String controlLamp(@PathVariable String deviceId, @PathVariable int status) {
        Map<String, Object> messagePayload = new HashMap<>();
        messagePayload.put("deviceId", deviceId);
        messagePayload.put("status", status);

        String json = JSON.toJSONString(messagePayload);
        String topic = "iot/lamp/server/status"; // Tópico para enviar comandos

        mqttMessageSender.sendToTopic(topic, json); // Usar el método simple del sender

        return "Comando enviado para device " + deviceId + ": status=" + status;
    }
}

3.4. Lógica de Registro de Estado del Disopsitivo

Tópico: iot/lamp/device/status
Mensaje (JSON): {"deviceId": "XXX", "status": 0} (0 para apagado, 1 para encendido)

Actualice ReceiverMessageHandler para manejar este tópico y llame al servicio correspondiente:

// Dentro de ReceiverMessageHandler, en handleMessage:
else if ("iot/lamp/device/status".equals(receivedTopic)) {
    tbLampStatusService.saveDeviceStatus(payload);
}

Implemente la lógica en TbLampStatusServiceImpl:

import com.alibaba.fastjson.JSON;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.example.demo.entity.TbLampStatus; // Entidad de estado de bombilla
import com.example.demo.mapper.TbLampStatusMapper; // Mapper de estado de bombilla
import com.example.demo.service.TbLampStatusService;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.Map;

@Service
public class TbLampStatusServiceImpl extends ServiceImpl<tblampstatusmapper tblampstatus=""> implements TbLampStatusService {

    @Override
    public void saveDeviceStatus(String json) {
        Map<string object=""> map = JSON.parseObject(json, Map.class);
        String deviceId = map.get("deviceId").toString();
        Integer status = Integer.parseInt(map.get("status").toString());

        TbLampStatus lampStatus = new TbLampStatus();
        lampStatus.setDeviceid(deviceId);
        lampStatus.setStatus(status);
        lampStatus.setCreateTime(new Date());
        save(lampStatus);
        System.out.println("Estado de dispositivo guardado: " + deviceId + ", Estado: " + status);
    }
}
</string></tblampstatusmapper>

Etiquetas: MQTT vue.js JavaScript java Spring Boot

Publicado el 9-9 03:30