Implémentation de la communication MQTT en Java

MQTT (Message Queuing Telemetry Transport) est un protocole de messagerie léger basé sur le modèle de publication-abonnement, conçu pour les environnements aux ressources limitées. Sa simplicité et son efficacité le rendent idéal pour la communication M2M (Machine-to-Machine) et l'IoT (Internet of Things).

Architecture Publication-Abonnement

Dans MQTT, les clients communiquent via un courtier (broker). Les clients peuvent publier des messages sur des "sujets" (topics) sans connaître les abonnés. Les clients qui s'abonnent à ces sujets reçoivent les messages pertinents. Ce découplage permet une grande flexibilité.

Outils et Bibliothèques

  • Serveur MQTT Broker : EMQX ou Mosquitto sont des options courantes pour héberger le courtier MQTT.
  • Clients de Test : MQTTX et MQTT.fx sont des outils graphiques pratiques pour tester la connectivité et échanger des messages.
  • Bibliothèque Client Java : La bibliothèque Paho d'Eclipse fournit une API Java robuste pour interagir avec les courtiers MQTT.

Intégration avec Spring Boot

L'intégration de MQTT dans une application Spring Boot peut être simplifiée en utilisant le module spring-integration-mqtt.

Dépendance Maven

<dependency>
   <groupId>org.springframework.integration</groupId>
   <artifactId>spring-integration-mqtt</artifactId>
</dependency>

Configuration MQTT (MqttConfig.java)

Cette classe configure les usines de clients MQTT et les adaptateurs pour la réception et l'envoi de messages.

import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.core.MessageProducer;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;

@Configuration
public class MqttConfig {

   // Configuration pour la réception de messages
   @Bean
   public MqttPahoClientFactory mqttClientFactory() {
       DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
       MqttConnectOptions options = new MqttConnectOptions();
       // Configuration des URI du serveur MQTT
       options.setServerURIs(new String[]{"tcp://localhost:1883"}); // Remplacez par l'adresse de votre broker
       factory.setConnectionOptions(options);
       return factory;
   }

   // Canal d'entrée pour les messages MQTT
   @Bean
   public MessageChannel mqttInputChannel() {
       return new DirectChannel();
   }

   // Adaptateur pour recevoir des messages MQTT
   @Bean
   public MessageProducer inbound() {
       // Le client ID "subscriber" est utilisé pour la connexion au broker
       MqttPahoMessageDrivenChannelAdapter adapter =
               new MqttPahoMessageDrivenChannelAdapter("subscriber", mqttClientFactory(), "sensor/data", "device/status");
       adapter.setCompletionTimeout(5000);
       adapter.setConverter(new DefaultPahoMessageConverter());
       adapter.setQos(1); // Niveau de Qualité de Service
       adapter.setOutputChannel(mqttInputChannel());
       return adapter;
   }

   // Gestionnaire pour traiter les messages reçus
   @Bean
   @ServiceActivator(inputChannel = "mqttInputChannel")
   public MessageHandler handler() {
       return message -> {
           String payload = new String((byte[]) message.getPayload()); // Assumant que le payload est en bytes
           String topic = (String) message.getHeaders().get("mqtt_receivedTopic");

           System.out.println("Message reçu sur le sujet : " + topic);
           System.out.println("Contenu : " + payload);

           // Logique de traitement des messages basée sur le sujet
           if (topic.startsWith("sensor/")) {
               System.out.println("Traitement des données du capteur...");
           } else if (topic.equals("device/status")) {
               System.out.println("Traitement du statut du périphérique...");
           }
       };
   }

   // Configuration pour l'envoi de messages
   // Canal de sortie pour les messages MQTT
   @Bean
   public MessageChannel mqttOutboundChannel() {
       return new DirectChannel();
   }

   // Gestionnaire pour envoyer des messages MQTT
   @Bean
   @ServiceActivator(inputChannel = "mqttOutboundChannel")
   public MessageHandler outbound() {
       // Le client ID "publisher" est utilisé pour la connexion au broker
       MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler("publisher", mqttClientFactory());
       messageHandler.setAsync(true); // Envoi asynchrone
       messageHandler.setDefaultTopic("command/execute"); // Sujet par défaut pour l'envoi
       messageHandler.setDefaultQos(1);
       messageHandler.setConverter(new DefaultPahoMessageConverter());
       return messageHandler;
   }
}

Interface Gateway (MqttGateway.java)

Cette interface définit des méthodes pour envoyer des messages via MQTT.

import org.springframework.integration.annotation.MessagingGateway;
import org.springframework.integration.mqtt.support.MqttHeaders;
import org.springframework.messaging.handler.annotation.Header;

@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MqttGateway {
   void sendToMqtt(String payload); // Envoi au sujet par défaut
   void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, String payload); // Envoi à un sujet spécifique
   void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) int qos, String payload); // Envoi avec QoS spécifié
   void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) int qos, byte[] payload); // Envoi en tant que tableau d'octets
}

Exemple de Contrôleur REST pour l'Envoi

Ce contrôleur permet d'envoyer des messages MQTT via une requête HTTP.

import com.example.mqtt.model.MessageRequest; // Assurez-vous que ce modèle existe
import com.example.mqtt.mqtt.MqttGateway;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;

import javax.annotation.Resource;

@RestController
public class MqttController {

   @Resource
   private MqttGateway mqttGateway;

   @PostMapping("/publish")
   public ResponseEntity<String> publishMessage(@RequestBody MessageRequest request) {
       String topic = request.getTopic();
       String message = request.getMessage();
       int qos = request.getQos();

       // Envoi du message en utilisant la gateway MQTT
       mqttGateway.sendToMqtt(topic, qos, message);

       String response = String.format("Message '%s' publié sur le sujet '%s' avec QoS %d.", message, topic, qos);
       return ResponseEntity.ok(response);
   }
}

Modèle de requête exemple (MessageRequest.java):

package com.example.mqtt.model;

public class MessageRequest {
   private String topic;
   private String message;
   private int qos;

   // Getters et Setters
   public String getTopic() { return topic; }
   public void setTopic(String topic) { this.topic = topic; }
   public String getMessage() { return message; }
   public void setMessage(String message) { this.message = message; }
   public int getQos() { return qos; }
   public void setQos(int qos) { this.qos = qos; }
}

Étiquettes: MQTT Java Spring Boot pub/sub iot

Publié le 15 septembre à 00h59