Apache RocketMQ offre des mécanismes puissants pour organiser et acheminer les messages via les concepts de topic et de tag. Le topic est essentiel pour déterminer la file d'attente à laquelle un message est acheminé, tandis que le tag sert de métadonnée secondaire pour un filtrage plus fin lors de la consommation.
Filtrage par Tag
Lors de l'envoi d'un message, le topic dirige le message vers une file d'attente spécifique. Le tag, quant à lui, n'influence pas la sélection de la file d'attente. Il est stocké en tant qu'attribut du message au sein de cette file d'attente. Au moment de la consommation, le broker utilise le tag pour acheminer les messages vers les consommateurs appropriés.
Exemple de producteur utilisant des tags
Ce producteur envoie deux messages sur le sujet tagTopic, chacun avec un tag distinct : vip1 et vip2.
// Initialisation du producteur
DefaultMQProducer producer = new DefaultMQProducer("tag-producer-group");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
// Création et envoi du premier message avec tag "vip1"
Message messageVip1 = new Message("tagTopic", "vip1", "Message VIP 1".getBytes());
producer.send(messageVip1);
// Création et envoi du second message avec tag "vip2"
Message messageVip2 = new Message("tagTopic", "vip2", "Message VIP 2".getBytes());
producer.send(messageVip2);
// Arrêt du producteur
producer.shutdown();
Exemple de consommateur 1 : récepsion de messages spécifiques
Ce consommateur est configuré pour recevoir uniquement les messages dont le tag est vip1.
// Initialisation du consommateur
DefaultMQPushConsumer consumerA = new DefaultMQPushConsumer("tag-consumer-group-a");
consumerA.setNamesrvAddr("127.0.0.1:9876");
// Abonnement au sujet "tagTopic" avec le filtre "vip1"
consumerA.subscribe("tagTopic", "vip1");
// Enregistrement du listener pour le traitement des messages
consumerA.registerMessageListener((List<MessageExt> msgs, ConsumeConcurrentlyContext context) -> {
System.out.println("Consommateur A reçoit : " + new String(msgs.get(0).getBody()));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
// Démarrage du consommateur
consumerA.start();
Exemple de consommateur 2 : réception de plusieurs tags
Ce consommateur peut recevoir des messages correspondant à plusieurs tags, ici vip1 ou vip2, en utilisant l'opérateur logique ||.
// Initialisation du consommateur
DefaultMQPushConsumer consumerB = new DefaultMQPushConsumer("tag-consumer-group-b");
consumerB.setNamesrvAddr("127.0.0.1:9876");
// Abonnement au sujet "tagTopic" avec les filtres "vip1" OU "vip2"
consumerB.subscribe("tagTopic", "vip1 || vip2");
// Enregistrement du listener pour le traitement des messages
consumerB.registerMessageListener((List<MessageExt> msgs, ConsumeConcurrentlyContext context) -> {
System.out.println("Consommateur B reçoit : " + new String(msgs.get(0).getBody()));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
// Démarrage du consommateur
consumerB.start();
Utilisation des Clés (Keys)
Les clés (key) offrent une autre dimension pour l'identification et la recherche de messages. Chaque message peut être associé à une ou plusieurs clés uniques. Bien qu'elles ne soient pas utilisées pour le routage initial comme les tags, les clés sont indexées par le broker et permettent des requêtes efficaces sur des messages spécifiques.
Exemple de producteur utilisant une clé
Ce producteur envoie un message sur keyTopic, avec le tag vip1 et une clé générée aléatoirement.
// Initialisation du producteur
DefaultMQProducer producerKey = new DefaultMQProducer("key-producer-group");
producerKey.setNamesrvAddr("127.0.0.1:9876");
producerKey.start();
// Génération d'une clé unique
String uniqueKey = UUID.randomUUID().toString();
// Création du message avec topic, tag, clé et corps du message
Message messageWithKey = new Message("keyTopic", "vip1", uniqueKey, "Message avec une clé unique".getBytes());
producerKey.send(messageWithKey);
// Arrêt du producteur
producerKey.shutdown();
Exemple de consommateur récupérant la clé du message
Ce consommateur s'abonne à keyTopic sans filtre spécifique (*) et affiche le corps ainsi que la clé du message reçu.
// Initialisation du consommateur
DefaultMQPushConsumer consumerKey = new DefaultMQPushConsumer("key-consumer-group");
consumerKey.setNamesrvAddr("127.0.0.1:9876");
// Abonnement à "keyTopic" pour tous les tags
consumerKey.subscribe("keyTopic", "*");
// Enregistrement du listener
consumerKey.registerMessageListener((List<MessageExt> msgs, ConsumeConcurrentlyContext context) -> {
MessageExt messageExt = msgs.get(0);
System.out.println("Corps du message : " + new String(messageExt.getBody()));
System.out.println("Clé(s) du message : " + messageExt.getKeys()); // Récupération de la clé
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
// Démarrage du consommateur
consumerKey.start();