Desarrollo Secundario de Protocolos Personalizados en Jetlinks

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:

  1. Modificación del código fuente.
  2. 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.

Etiquetas: jetlinks MQTT protocolos personalizados IoT java

Publicado el 9-20 05:52