Les modes de fonctionnement de RabbitMQ

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

  1. 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();
    }
}



  1. 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
  1. 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

  1. 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);
    }
}



  1. 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);

    }
}



Étiquettes: rabbitmq files d'attente de messages échangeurs routage mode thématique

Publié le 17 septembre à 04h40