Este artículo detalla cómo implementar el procesamiento de protocolos personalizados para un sistema de Internet de las Cosas (IoT) construido sobre Jetlinks y EMQX, permitiendo la integración directa de dispositivos vía MQTT y la conexión a través de un broker MQTT.
Implementación del Manejo de Protocolos Personalizados
Existen dos enfoques principales para manejar protocolos MQTT personalizados y la comuincación con sistemas de negocio:
- Modificación del código fuente.
- Desarrollo secundario sobre los protocolos oficiales.
Enfoque 1: Modificación del Código Fuente
Este método implica suscribirse a los topics de mensajes MQTT y utilizar HTTP o colas de mensajes (como RabbitMQ) para enviar los datos a los sistemas de negocio para su procesamiento.
@Subscribe("/device/*/*/message/property/report")
public Mono<Void> handlePropertyReport(TopicPayload payload) {
// Enviar a RabbitMQ para procesamiento asíncrono
rabbitMqProducer.send(payload, MessageType.PROPERTIES, RabbitConfig.RabbitConstants.QUEUE_QUEUE_PROPERTIES);
// Enviar directamente usando un productor RabbitMQ
rabbitMqSender.send(payload, MessageType.PROPERTIES);
return Mono.empty();
}
@Subscribe("/device/*/*/message/event/*")
public Mono<Void> handleDeviceEvent(TopicPayload payload) {
// Enviar a RabbitMQ para procesamiento asíncrono de eventos
rabbitMqProducer.send(payload, MessageType.EVENT, RabbitConfig.RabbitConstants.QUEUE_QUEUE_EVENT);
// Enviar directamente usando un productor RabbitMQ para eventos
rabbitMqSender.send(payload, MessageType.EVENT);
return Mono.empty();
}
Enfoque 2: Desarrollo Secundario de Protocolos
Personalización de Cabeceras de Mensajes en TopicMessageCodec
Se pueden definir manejadores personalizados para interactuar con tipos específicos de dispositivos o mensajes.
// Manejador para respuestas de propiedades de "Máquinas de Información" (XXJ)
xxjCallbackProperty("/*/properties/callback/xxj",
CapexMessage.class,
route -> route
.upstream(false) // No enviar upstream a Jetlinks
.downstream(true) // Permitir enviar downstream desde Jetlinks
.group("Respuestas de Dispositivo")
.description("Respuesta activa de información del dispositivo")
.example("{\"properties\":{\"IDPropiedad\":\"ValorPropiedad\"}}")),
// Manejador para respuestas de propiedades de "Torniquetes" (HKMJ)
hkmjCallbackProperty("/*/properties/callback/hkmj",
CapexMessage.class,
route -> route
.upstream(false)
.downstream(true)
.group("Respuestas de Dispositivo")
.description("Respuesta activa de información del dispositivo")
.example("{\"properties\":{\"IDPropiedad\":\"ValorPropiedad\"}}")),
Intercepción y Análisis de Mensajes MQTT en MqttDeviceMessageCodec
Se puede sobreescribir el método decode para interceptar y procesar mensajes MQTT entrantes según la lógica personalizada. El código de ejemplo para este desarrollo secundario está disponible en GitHub.
@Nonnull
@Override
public Flux<DeviceMessage> decode(@Nonnull MessageDecodeContext context) {
MqttMessage mqttMessage = (MqttMessage) context.getMessage();
byte[] messageBytes = mqttMessage.payloadAsBytes();
// TODO: Implementar lógica personalizada para el reporte de dispositivos (XXJ, HKMJ)
String messageContent = new String(messageBytes);
// Procesamiento personalizado del mensaje recibido del equipo
processCustomMessageFromEquipment(messageContent, context);
// Intentar decodificar usando el manejador de topic estándar
return TopicMessageCodec
.decode(mapper, TopicMessageCodec.removeProductPath(mqttMessage.getTopic()), messageBytes)
// Si la decodificación directa falla, podría ser otra funcionalidad del dispositivo
.switchIfEmpty(FunctionalTopicHandlers
.handle(context.getDevice(),
mqttMessage.getTopic().split("/"),
messageBytes,
mapper,
reply -> performReply(context, reply)));
}
El repositorio para este desarrollo secundario de protocolo se encuentra en: jetlinks-protocol-gt.