今天过节?我不管,我就是要学RabbitMQ Work Queues(Java版)

Posted 你,好

tags:

篇首语:本文由小常识网(cha138.com)小编为大家整理,主要介绍了今天过节?我不管,我就是要学RabbitMQ Work Queues(Java版)相关的知识,希望对你有一定的参考价值。

虽然看我博客的人很少,但是还是祝大家中国(中秋+国庆)节快乐,学业有成,offer拿到手软哦。
今天带来的是关于RabbitMQ的Work Queues的学习笔记。

Work Queues 也叫 Task Queues。当生产者将大量任务发送至消息队列时,消息队列默认会按顺序将任务分发给每个消费者,这样多个消费者同时执行任务,极大的提高了效率。
这里我们模拟一个生产者,两个消费者,来进行演示。

public class Provider 
    public static void main(String[] args) 
        try 
            Connection connection = Util.getConnection(); // 获取连接
            Channel channel = connection.createChannel(); // 创建频道
            channel.queueDeclare("work", false, false, false, null); // 声明队列
            for (int i = 0; i < 30; ++i) 
                channel.basicPublish("", "work", null, ("任务:" + i).getBytes()); // 发布任务
            
            Util.close(channel, connection); // 关闭连接
         catch (Exception e) 
            e.printStackTrace();
        
    

public class Consumer1 
    public static void main(String[] args) 
        try 
            Connection connection = Util.getConnection(); // 获取连接
            Channel channel = connection.createChannel(); // 创建频道
            channel.queueDeclare("work",false, false, false, null); // 声明队列
            // 处理任务
            channel.basicConsume("work", true, new DefaultConsumer(channel)
                @Override
                public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException 
                    System.out.println("消费者-1 " + new String(body)); 
                    try 
                        Thread.sleep(2000);
                     catch (InterruptedException e) 
                        e.printStackTrace();
                    
                
            );
         catch (Exception e) 
            e.printStackTrace();
        
    


public class Consumer2 
    public static void main(String[] args) 
        try 
            Connection connection = Util.getConnection();
            Channel channel = connection.createChannel();
            channel.queueDeclare("work",false, false, false, null);
            channel.basicConsume("work", true, new DefaultConsumer(channel)
                @Override
                public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException 
                    System.out.println("消费者-2 " + new String(body));
                    try 
                        Thread.sleep(1000);
                     catch (InterruptedException e) 
                        e.printStackTrace();
                    
                
            );
         catch (Exception e) 
            e.printStackTrace();
        
    

其中,我们用Thread.sleep()模拟消费者处理任务消耗的时间。Consumer1每处理一个任务需要2秒,Consumer2则只需要1秒。而连接RabbitMQ的代码,复用性高,我直接封装成了工具类。

public class Util 
    private static final ConnectionFactory connectionFactory;
    static 
        connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("xx.xx.xx.xx");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/test");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
    

    public static Connection getConnection() throws IOException, TimeoutException 
        return connectionFactory.newConnection();
    

    public static void close(Channel channel, Connection con) throws IOException, TimeoutException 
        if (channel != null) channel.close();
        if (con != null) con.close();
    

运行,得到测验成果

可以看到,任务平分给了两个消费者,但是Consumer1处理完的时间是Consumer2的两倍。所以在任务量巨大的情况下,这种平分的方式将会造成巨大的资源浪费,可能Consumer2早早的处理完任务,但Consumer1还差很远。所以为了充分利用每个消费者资源,做到能者多劳,我们需要做些改进。
首先,需要将自动确认(autoAck)设为false,因为我们需要让每个任务处理完毕后再确认。
basicConsume(String queue, boolean autoAck, Consumer callback)
其次,手动确认任务。
basicAck(long deliveryTag, boolean multiple)
最后,把消息队列分发给消费者的任务数量设为1。这样消息队列就不会一股脑的把任务分完,而是会观察哪个消费者有空闲,再进行分发。
basicQos(int prefetchCount)
完整代码(以Consumer1为例,更改的地方Consumer2同Consumer1)

public class Consumer1 
    public static void main(String[] args) 
        try 
            Connection connection = Util.getConnection();
            Channel channel = connection.createChannel();
            channel.basicQos(1); // 设置分发任务数量为1
            channel.queueDeclare("work",false, false, false, null);
            boolean autoAck = false; // 自动确认为false
            channel.basicConsume("work", autoAck, new DefaultConsumer(channel)
                @Override
                public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException 
                    System.out.println("消费者-1 " + new String(body));
                    try 
                        Thread.sleep(2000);
                     catch (InterruptedException e) 
                        e.printStackTrace();
                    
                    channel.basicAck(envelope.getDeliveryTag(), false); // 处理完任务手动确认
                
            );
         catch (Exception e) 
            e.printStackTrace();
        
    

结果,明显看到,能者多劳。
此程序同样解决了当其中一台消费者停机时,数据丢失的问题。道理也很简单,消费者停机后,消息队列就不再对他分配任务,而是会将任务分配给正常运行的消费者。
但是,如果RabbitMQ重启了,或者消费者与队列的连接、频道断开了,目前的程序就无法避免数据丢失了。所以我们需要将存储任务的队列设置为可持久化队列。
操作很简单。
首先,在声明队列时将durable参数设为true。
queueDeclare(String queue, boolean durable, boolean exclusive, boolean autoDelete, Map<String, Object> arguments)
之后在发送消息时将消息参数设为MessageProperties.PERSISTENT_TEXT_PLAIN即可。
因为之前的work队列已设置durable为false,所以我们新建一个名为work-durable队列做演示。

public class Provider 
    public static void main(String[] args) 
        try 
            Connection connection = Util.getConnection();
            Channel channel = connection.createChannel();
            channel.queueDeclare("work-durable", true, false, false, null); // 持久化设为true
            for (int i = 0; i < 30; ++i) 
                channel.basicPublish("", "work-durable", MessageProperties.PERSISTENT_TEXT_PLAIN, ("任务:" + i).getBytes()); // 消息参数设为PERSISTENT_TEXT_PLAIN
            
            Util.close(channel, connection);
         catch (Exception e) 
            e.printStackTrace();
        
    

消费者就不再赘述。
这样,我们的数据在一般情况下就不会轻易丢失了。

以上是关于今天过节?我不管,我就是要学RabbitMQ Work Queues(Java版)的主要内容,如果未能解决你的问题,请参考以下文章

爬虫入门及HTTP协议的讲解

Golang调用Rabbitmq消息队列和封装

捋明白 RabbitMQ 中的权限系统

RabbitMQ在特来电的深度应用

来去匆匆

阿里程序员,30 岁之前你要学着开始规划自己的职业发展了