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; }
}