Démo RabbitMQ + Intégration Spring/SpringBoot avec RabbitMQ
(Lien du dépôt GitHub)[https://github.com/2537422279/RabbitMQDemo]
Envoi de messages basé sur des files d'attente
- Bonjour le monde ! (File d'attente de messages de base) ======================================
Implémentation du code
consommateur
package fr.rabbitmq.example.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/03 16:33
*/
public class Receveur_Bonjour {
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setPort(5672);
factory.setVirtualHost("/");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
Consumer consommateur = new DefaultConsumer(channel){
/*
Méthode de rappel, exécutée automatiquement lors de la réception d'un message.
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("consumerTag: " + consumerTag);
System.out.println("Exchange: " + envelope.getExchange());
System.out.println("RoutingKey: " + envelope.getRoutingKey());
System.out.println("Properties: " + properties);
System.out.println("body: " + new String(body));
}
};
channel.basicConsume("bonjour_monde",true,consommateur);
}
}
producteur
package fr.rabbitmq.producteur;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @description Envoi de messages
* @author Auteur
* @version V1.0
* @date 2022/10/03 15:08
*/
public class Producteur_Bonjour {
public static void main(String[] args) throws IOException, TimeoutException {
// Création de la fabrique de connexion
ConnectionFactory factory = new ConnectionFactory();
// Configuration des paramètres,
factory.setHost("localhost"); // Adresse IP
factory.setPort(5672); // Port, valeur par défaut 5672
factory.setVirtualHost("/"); // Configuration de la machine virtuelle, valeur par défaut "/"
factory.setUsername("utilisateur");
factory.setPassword("motdepasse");
// Création de la connexion
final Connection connection = factory.newConnection();
// Création du canal
final Channel channel = connection.createChannel();
// Création de la file d'attente de messages
/*
queueDeclare(String queue, boolean durable, boolean exclusive, boolean autoDelete, Map<String, Object> arguments) throws IOException {
*/
channel.queueDeclare("bonjour_monde", true, false, false, null);
// Envoi du message
String corps = "bonjour rabbitmq~~~~";
channel.basicPublish("","bonjour_monde",null,corps.getBytes());
// Libération des ressources
channel.close();
connection.close();
}
}
- File d'attente de travail ======================
producteur
package fr.rabbitmq.producteur;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/03 17:36
*/
public class Producteur_FilesTravail {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setPort(5672);
factory.setHost("localhost");
factory.setVirtualHost("/");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
// Création de la file d'attente de messages
channel.queueDeclare("files_travail",true,false,false,null);
for (int i = 1; i <= 10; i++) {
String corps = i + "bonjour rabbitmq~~~";
channel.basicPublish("","files_travail",null,corps.getBytes());
}
}
}
consommateur1
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/03 17:47
*/
public class Receveur_FilesTravail1 {
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setPort(5672);
factory.setVirtualHost("/");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
Consumer consommateur = new DefaultConsumer(channel){
/*
Méthode de rappel, exécutée automatiquement lors de la réception d'un message.
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
// System.out.println("consumerTag: " + consumerTag);
// System.out.println("Exchange: " + envelope.getExchange());
// System.out.println("RoutingKey: " + envelope.getRoutingKey());
// System.out.println("Properties: " + properties);
System.out.println("body: " + new String(body));
}
};
channel.basicConsume("files_travail",true,consommateur);
}
}
consommateur2
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/03 17:47
*/
public class Receveur_FilesTravail2 {
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setPort(5672);
factory.setVirtualHost("/");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
Consumer consommateur = new DefaultConsumer(channel){
/*
Méthode de rappel, exécutée automatiquement lors de la réception d'un message.
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
// System.out.println("consumerTag: " + consumerTag);
// System.out.println("Exchange: " + envelope.getExchange());
// System.out.println("RoutingKey: " + envelope.getRoutingKey());
// System.out.println("Properties: " + properties);
System.out.println("body: " + new String(body));
}
};
channel.basicConsume("files_travail",true,consommateur);
}
}
Lancez d'abord les deux consommateurs, puis le producteur. L'effet est le suivant
Résumé
- Les consommateurs sont en compétition pour le même message
- Les consommateurs alternent pour lire les messages de la file d'attente
- Mode de fontcionnement Pub/Sub =======================
Introduction d'un échangeur
Producteur_PubSub
package fr.rabbitmq.producteur;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 16:05
*/
public class Producteur_PubSub {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("utilisateur");
factory.setPassword("motdepasse");
factory.setVirtualHost("/");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
// Création de l'échangeur
// String exchange, BuiltinExchangeType type, boolean durable, boolean autoDelete, boolean internal, Map<String, Object> arguments
// exchange:Nom de l'échangeur
// type:
String nomEchangeur = "test_fanout";
channel.exchangeDeclare(nomEchangeur, BuiltinExchangeType.FANOUT,true,false,false,null);
// Création des files d'attente
String file1Nom = "test_fanout_file1";
String file2Nom = "test_fanout_file2";
channel.queueDeclare(file1Nom,true,false,false,null);
channel.queueDeclare(file2Nom,true,false,false,null);
// Liaison des files d'attente et de l'échangeur
channel.queueBind(file1Nom,nomEchangeur,"");
channel.queueBind(file2Nom,nomEchangeur,"");
// Envoi du message
String corps = "Information de journal : Zhang a appelé la méthode fanout.... Niveau de journalisation : info...";
channel.basicPublish(nomEchangeur,"",null,corps.getBytes());
channel.close();
connection.close();
}
}
Receveur_PubSub1
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 16:25
*/
public class Receveur_PubSub1 {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setVirtualHost("/");
factory.setPort(5672);
factory.setUsername("utilisateur");
factory.setPassword("motdepasse");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String file1Nom = "test_fanout_file1";
String file2Nom = "test_fanout_file2";
final Consumer consommateur = new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("body: " + new String(body));
System.out.println("Imprimer les informations de journal sur la console....");
}
};
channel.basicConsume(file1Nom,true,consommateur);
}
}
Receveur_PubSub2
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 16:25
*/
public class Receveur_PubSub2 {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setVirtualHost("/");
factory.setPort(5672);
factory.setUsername("utilisateur");
factory.setPassword("motdepasse");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String file1Nom = "test_fanout_file1";
String file2Nom = "test_fanout_file2";
final Consumer consommateur = new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("body: " + new String(body));
System.out.println("Enregistrer les informations de journal dans la base de données....");
}
};
channel.basicConsume(file2Nom,true,consommateur);
}
}
Affichage des résultats
- Mode de routage ================
Producteur_Routage
package fr.rabbitmq.producteur;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 16:45
*/
public class Producteur_Routage {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("utilisateur");
factory.setPassword("motdepasse");
factory.setVirtualHost("/");
factory.setPort(5672);
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String nomEchangeur="test_direct";
channel.exchangeDeclare(nomEchangeur, BuiltinExchangeType.DIRECT,true,false,false,null);
String file1Nom = "test_direct_file1";
String file2Nom = "test_direct_file2";
channel.queueDeclare(file1Nom,true,false,false,null);
channel.queueDeclare(file2Nom,true,false,false,null);
channel.queueBind(file1Nom,nomEchangeur,"erreur");
channel.queueBind(file2Nom,nomEchangeur,"erreur");
channel.queueBind(file2Nom,nomEchangeur,"info");
channel.queueBind(file2Nom,nomEchangeur,"avertissement");
String corps = "Information de journal : Zhang a appelé la méthode de routage... Niveau de journalisation : info...";
channel.basicPublish(nomEchangeur,"avertissement",null,corps.getBytes());
channel.close();
connection.close();
}
}
Receveur_Routage1
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 17:26
*/
public class Receveur_Routage1 {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setPort(5672);
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setVirtualHost("/");
factory.setHost("localhost");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String file1Nom = "test_direct_file1";
String file2Nom = "test_direct_file2";
Consumer consommateur = new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("body: " + new String(body));
System.out.println("Imprimer les informations de journal sur la console....");
}
};
channel.basicConsume(file2Nom,true,consommateur);
}
}
Receveur_Routage2
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 17:26
*/
public class Receveur_Routage2 {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setPort(5672);
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setVirtualHost("/");
factory.setHost("localhost");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String file1Nom = "test_direct_file1";
String file2Nom = "test_direct_file2";
Consumer consommateur = new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("body: " + new String(body));
System.out.println("Enregistrer les informations de journal dans la base de données....");
}
};
channel.basicConsume(file1Nom,true,consommateur);
}
}
- Mode Thématique ===============
Producteur_Theme
package fr.rabbitmq.producteur;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 19:09
*/
public class Producteur_Theme {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("utilisateur");
factory.setPassword("motdepasse");
factory.setVirtualHost("/");
factory.setPort(5672);
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String nomEchangeur="test_theme";
channel.exchangeDeclare(nomEchangeur, BuiltinExchangeType.TOPIC,true,false,false,null);
String file1Nom = "test_theme_file1";
String file2Nom = "test_theme_file2";
channel.queueDeclare(file1Nom,true,false,false,null);
channel.queueDeclare(file2Nom,true,false,false,null);
channel.queueBind(file1Nom,nomEchangeur,"#.erreur");
channel.queueBind(file1Nom,nomEchangeur,"commande.*");
channel.queueBind(file2Nom,nomEchangeur,"*.*");
String corps = "Je suis un message...";
channel.basicPublish(nomEchangeur,"marchandise.erreur",null,corps.getBytes());
}
}
Receveur_Theme1
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 19:26
*/
public class Receveur_Theme1 {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setPort(5672);
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setVirtualHost("/");
factory.setHost("localhost");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String file1Nom = "test_theme_file1";
String file2Nom = "test_theme_file2";
Consumer consommateur = new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("body: " + new String(body));
System.out.println("Enregistrer les informations de journal dans la base de données....");
}
};
channel.basicConsume(file1Nom,true,consommateur);
}
}
Receveur_Theme2
package fr.rabbitmq.receveur;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @author Auteur
* @version V1.0
* @date 2022/10/04 19:26
*/
public class Receveur_Theme2 {
public static void main(String[] args) throws IOException, TimeoutException {
final ConnectionFactory factory = new ConnectionFactory();
factory.setPort(5672);
factory.setPassword("motdepasse");
factory.setUsername("utilisateur");
factory.setVirtualHost("/");
factory.setHost("localhost");
final Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();
String file1Nom = "test_theme_file1";
String file2Nom = "test_theme_file2";
Consumer consommateur = new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("body: " + new String(body));
System.out.println("Imprimer les informations de journal sur la console....");
}
};
channel.basicConsume(file2Nom,true,consommateur);
}
}