RabbitMQ->How2j教程

1.安装erlang与RabbitMQ,能正常访问到rabbitMQ的界面即表示ok,默认用户名及密码都是guest


2.一些相关概念介绍:

(1).AMQP通信协议:可以处理更为复杂的协议。像http,socket,ftp,都是通信协议。
(2).消息路由过程:RabbitMQ拿到消息后,会先给到交换机(Exchange),然后交换机再根据预先设定不同的绑定策略(Bindings),来确定要发给哪个队列。
(3)模式:RabbitMQ提供了四种Exchange模式:fanout,direct,topic,header。主要是前三种:fanot:广播模式,发给所有队列。Direct模式:只发给指定的队列,其他的队列不进行发送。topic模式:通过指定一个主题,会根据不同的队列分发不同的内容。

3.代码示例
根据站长的例子,这里我们仅仅做简单的Demo,如果有更详细的需要去后续自行了解,这里将能跑通的demo记录一下,以后使用的时候可以用到项目中来。

fanout

pom.xml文件

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>
  <groupId>com.btzh</groupId>
  <artifactId>rabbitmq-test</artifactId>
  <version>0.0.1-SNAPSHOT</version>
  <name>rabbitmq</name>
  <description>rabbitmq</description>
   <dependencies>
    <dependency>
            <groupId>com.rabbitmq</groupId>
            <artifactId>amqp-client</artifactId>
            <version>3.6.5</version>
     </dependency>
    <dependency>
            <groupId>com.rabbitmq</groupId>
            <artifactId>amqp-client</artifactId>
            <version>5.6.0</version>
</dependency>     
   </dependencies>
</project>
TestProducer类:消息生产者:
package com.btzh.rabbitmq_test;

import java.io.IOException;
import java.util.concurrent.TimeoutException;
 
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
 
/**
 * 消息生成者
 */
public class TestProducer {
    public final static String EXCHANGE_NAME="fanout_exchange";
 
    public static void main(String[] args) throws IOException, TimeoutException {
         
        //创建连接工厂
        ConnectionFactory factory = new ConnectionFactory();
        //设置RabbitMQ相关信息
        factory.setHost("localhost");
        //创建一个新的连接
        Connection connection = factory.newConnection();
        //创建一个通道
        Channel channel = connection.createChannel();
         
        channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
         
        for (int i = 0; i < 100; i++) {
            String message = "direct 消息 " +i;
            //发送消息到队列中
            channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes("UTF-8"));
            System.out.println("发送消息: " + message);
             
        }
        //关闭通道和连接
        channel.close();
        connection.close();
    }
}
TestCustomer:消息消费者类:
package com.btzh.rabbitmq_test;

import java.io.IOException;
import java.util.concurrent.TimeoutException;
 
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
 
 
public class TestCustomer {
    public final static String EXCHANGE_NAME="fanout_exchange";     //交换器名称
 
    public static void main(String[] args) throws IOException, TimeoutException {
        //为当前消费者取随机名
        final String name = "consumer-"+ Math.random();   
         // 创建连接工厂
        ConnectionFactory factory = new ConnectionFactory();
        //设置RabbitMQ地址
        factory.setHost("localhost");
        //创建一个新的连接
        Connection connection = factory.newConnection();
        //创建一个通道
        Channel channel = connection.createChannel();
        //交换机声明(参数为:交换机名称;交换机类型)
        channel.exchangeDeclare(EXCHANGE_NAME,"fanout");
        //获取一个临时队列
        String queueName = channel.queueDeclare().getQueue();
        //队列与交换机绑定(参数为:队列名称;交换机名称;routingKey忽略)
        channel.queueBind(queueName,EXCHANGE_NAME,"");
         
        System.out.println(name +" 等待接受消息");
        //DefaultConsumer类实现了Consumer接口,通过传入一个频道,
        // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery
        Consumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope,
                                       AMQP.BasicProperties properties, byte[] body)
                    throws IOException {
                String message = new String(body, "UTF-8");
                System.out.println(name + " 接收到消息 '" + message + "'");
            }
        };
        //自动回复队列应答 -- RabbitMQ中的消息确认机制
        channel.basicConsume(queueName, true, consumer);
    }
}


direct

Provider:
package com.btzh.rabbitmq_test_direct;

/**
 * @author wzy
 *
 */
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
     
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
     
    /**
     * 消息生成者
     */
    public class TestDriectProducer {
        public final static String QUEUE_NAME="direct_queue";       //
     
        public static void main(String[] args) throws IOException, TimeoutException {
         
             
            //创建连接工厂
            ConnectionFactory factory = new ConnectionFactory();
            //设置RabbitMQ相关信息
            factory.setHost("localhost");
            //创建一个新的连接
            Connection connection = factory.newConnection();
            //创建一个通道
            Channel channel = connection.createChannel();
             
            for (int i = 0; i < 100; i++) {
                String message = "direct 消息 " +i;
                //发送消息到队列中
                channel.basicPublish("", QUEUE_NAME, null, message.getBytes("UTF-8"));
                System.out.println("发送消息: " + message);
                 
            }
            //关闭通道和连接
            channel.close();
            connection.close();
        }
    }
Consumer:
package com.btzh.rabbitmq_test_direct;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
 
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
 
 
public class TestDriectCustomer {
    private final static String QUEUE_NAME = "direct_queue";
 
    public static void main(String[] args) throws IOException, TimeoutException {
        //为当前消费者取随机名
        final String name = "consumer-"+ Math.random()*10;
         
        // 创建连接工厂
        ConnectionFactory factory = new ConnectionFactory();
        //设置RabbitMQ地址
        factory.setHost("localhost");
        //创建一个新的连接
        Connection connection = factory.newConnection();
        //创建一个通道
        Channel channel = connection.createChannel();
        //声明要关注的队列
        channel.queueDeclare(QUEUE_NAME, false, false, true, null);
        System.out.println(name +" 等待接受消息");
        //DefaultConsumer类实现了Consumer接口,通过传入一个频道,
        // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery
        Consumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope,
                                       AMQP.BasicProperties properties, byte[] body)
                    throws IOException {
                String message = new String(body, "UTF-8");
                System.out.println(name + " 接收到消息 '" + message + "'");
            }
        };
        //自动回复队列应答 -- RabbitMQ中的消息确认机制
        channel.basicConsume(QUEUE_NAME, true, consumer);
    }
}


Topic:

Provider:
package com.btzh.rabbitmq_test_topic;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

/**
 * @author wzy
 *
 */
public class TestProducer {
      public final static String EXCHANGE_NAME="topics_exchange";
      
        public static void main(String[] args) throws IOException, TimeoutException {
          
            //创建连接工厂
            ConnectionFactory factory = new ConnectionFactory();
            //设置RabbitMQ相关信息
            factory.setHost("localhost");
            //创建一个新的连接
            Connection connection = factory.newConnection();
            //创建一个通道
            Channel channel = connection.createChannel();
             
            channel.exchangeDeclare(EXCHANGE_NAME, "topic");
             
            String[] routing_keys = new String[] { "usa.news", "usa.weather",   
                    "europe.news", "europe.weather" };   
            String[] messages = new String[] { "美国新闻", "美国天气",   
                    "欧洲新闻", "欧洲天气" };   
             
            for (int i = 0; i < routing_keys.length; i++) {
                String routingKey = routing_keys[i];
                String message = messages[i];
                channel.basicPublish(EXCHANGE_NAME, routingKey, null, message   
                        .getBytes());   
                System.out.printf("发送消息到路由:%s, 内容是: %s%n ", routingKey,message);
                 
            }
     
            //关闭通道和连接
            channel.close();
            connection.close();
        }
}
接受*.news的队列:
package com.btzh.rabbitmq_test_topic;

import java.io.IOException;
import java.util.concurrent.TimeoutException;
 
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

 
public class TestCustomer4News {
    public final static String EXCHANGE_NAME="topics_exchange";
 
    public static void main(String[] args) throws IOException, TimeoutException {
        //为当前消费者取名称
        final String name = "consumer-news";
         
        // 创建连接工厂
        ConnectionFactory factory = new ConnectionFactory();
        //设置RabbitMQ地址
        factory.setHost("localhost");
        //创建一个新的连接
        Connection connection = factory.newConnection();
        //创建一个通道
        Channel channel = connection.createChannel();
        //交换机声明(参数为:交换机名称;交换机类型)
        channel.exchangeDeclare(EXCHANGE_NAME,"topic");
        //获取一个临时队列
        String queueName = channel.queueDeclare().getQueue();
        //接受 USA 信息
         
        channel.queueBind(queueName, EXCHANGE_NAME, "*.news");           
        System.out.println(name +" 等待接受消息");
        //DefaultConsumer类实现了Consumer接口,通过传入一个频道,
        // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery
        Consumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope,
                                       AMQP.BasicProperties properties, byte[] body)
                    throws IOException {
                String message = new String(body, "UTF-8");
                System.out.println(name + " 接收到消息 '" + message + "'");
            }
        };
        //自动回复队列应答 -- RabbitMQ中的消息确认机制
        channel.basicConsume(queueName, true, consumer);
    }
}
接收usa.*f的所有信息
package com.btzh.rabbitmq_test_topic;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
 
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

public class TestCustomer4USA {
    public final static String EXCHANGE_NAME="topics_exchange";
 
    public static void main(String[] args) throws IOException, TimeoutException {
        //为当前消费者取名称
        final String name = "consumer-usa";

        // 创建连接工厂
        ConnectionFactory factory = new ConnectionFactory();
        //设置RabbitMQ地址
        factory.setHost("localhost");
        //创建一个新的连接
        Connection connection = factory.newConnection();
        //创建一个通道
        Channel channel = connection.createChannel();
        //交换机声明(参数为:交换机名称;交换机类型)
        channel.exchangeDeclare(EXCHANGE_NAME,"topic");
        //获取一个临时队列
        String queueName = channel.queueDeclare().getQueue();
        //接受 USA 信息
         
        channel.queueBind(queueName, EXCHANGE_NAME, "usa.*");           
        System.out.println(name +" 等待接受消息");
        //DefaultConsumer类实现了Consumer接口,通过传入一个频道,
        // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery
        Consumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope,
                                       AMQP.BasicProperties properties, byte[] body)
                    throws IOException {
                String message = new String(body, "UTF-8");
                System.out.println(name + " 接收到消息 '" + message + "'");
            }
        };
        //自动回复队列应答 -- RabbitMQ中的消息确认机制
        channel.basicConsume(queueName, true, consumer);
    }
}

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容