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.
- 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>
- 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.
- 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
- Base de Datos: Cree las tablas
tb_lamp(para estado online/offline del dispositivo) ytb_lamp_status(para estado on/off de la bombilla). - Dependencias Adicionales: Añada
mybatis-plus-boot-startery el driver de MySQL a supom.xml. - Configuración de DataSource: Configure la conexión a la base de datos en
application.yml. - Generación de Código MyBatis-Plus: Use el generador de código para crear los Mappers y Entidades.
- Escaneo de Mappers: Añada la anotación
@MapperScana 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>