好得很程序员自学网

<tfoot draggable='sEl'></tfoot>

在RabbitMQ中实现Work queues工作队列模式

一、模式说明

Work Queues 与入门程序的简单模式相比,多了一个或一些消费端,多个消费端共同消费同一个队列中的消息。

应用场景 :对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。

二、代码

Work Queues 与入门程序的 简单模式 的代码是几乎一样的:可以完全复制,并复制多一个消费者进行多个消费者同时消费消息的测试。

①生产者

?

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

17

18

19

20

21

22

23

24

25

26

27

28

29

30

31

32

33

34

35

36

37

package com.itheima.rabbitmq.work;

import com.itheima.rabbitmq.util.ConnectionUtil;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class Producer {

     static final String QUEUE_NAME = "work_queue" ;

     public static void main(String[] args) throws Exception {

         //创建连接

         Connection connection = ConnectionUtil.getConnection();

         // 创建频道

         Channel channel = connection.createChannel();

         // 声明(创建)队列

         /**

          * 参数1:队列名称

          * 参数2:是否定义持久化队列

          * 参数3:是否独占本次连接

          * 参数4:是否在不使用的时候自动删除队列

          * 参数5:队列其它参数

         */

         channel.queueDeclare(QUEUE_NAME, true , false , false , null );

         for ( int i = 1 ; i <= 30 ; i++) {

             // 发送信息

             String message = "你好;小兔子!work模式--" + i;

             /**

              * 参数1:交换机名称,如果没有指定则使用默认Default Exchage

              * 参数2:路由key,简单模式可以传递队列名称

              * 参数3:消息其它属性

              * 参数4:消息内容

             */

             channel.basicPublish( "" , QUEUE_NAME, null , message.getBytes());

             System.out.println( "已发送消息:" + message);

         }

         // 关闭资源

         channel.close(); connection.close();

     }

}

②消费者1

?

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

17

18

19

20

21

22

23

24

25

26

27

28

29

30

31

32

33

34

35

36

37

38

39

40

41

42

43

44

45

46

47

48

49

50

51

52

53

54

55

56

57

58

package com.itheima.rabbitmq.work;

import com.itheima.rabbitmq.util.ConnectionUtil;

import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer1 {

     public static void main(String[] args) throws Exception {

         Connection connection = ConnectionUtil.getConnection();

         // 创建频道

         Channel channel = connection.createChannel();

         // 声明(创建)队列

         /**

          * 参数1:队列名称

          * 参数2:是否定义持久化队列

          * 参数3:是否独占本次连接

          * 参数4:是否在不使用的时候自动删除队列

          * 参数5:队列其它参数

         */

         channel.queueDeclare(Producer.QUEUE_NAME, true , false , false , null );

         //一次只能接收并处理一个消息

         channel.basicQos( 1 );

         //创建消费者;并设置消息处理

         DefaultConsumer consumer = new DefaultConsumer(channel){

             @Override

             /**

              * consumerTag 消息者标签,在channel.basicConsume时候可以指定

              * envelope 消息包的内容,可从中获取消息id,消息routingkey,交换机,消息和重传标志(收到消息失败后是否需要重新发送)

              * properties 属性信息

              * body 消息

             */

             public void handleDelivery(String consumerTag, Envelope envelope,

                     AMQP.BasicProperties properties, byte [] body) throws IOException {

                 try {

                     //路由key

                     System.out.println( "路由key为:" + envelope.getRoutingKey());

                     //交换机

                     System.out.println( "交换机为:" + envelope.getExchange());

                     //消息id

                     System.out.println( "消息id为:" + envelope.getDeliveryTag());

                     //收到的消息

                     System.out.println( "消费者1-接收到的消息为:" + new String(body, "utf-8" ));

                     Thread.sleep( 1000 );

                     //确认消息

                     channel.basicAck(envelope.getDeliveryTag(), false );

                 }

                 catch (InterruptedException e) {

                     e.printStackTrace();

                 }

             }

         };

         //监听消息

         /**

          * 参数1:队列名称

          * 参数2:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复会删除消息,设置为false则需要手动确认

          * 参数3:消息接收到后回调

         */

         channel.basicConsume(Producer.QUEUE_NAME, false , consumer);

     }

}

③消费者2

?

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

17

18

19

20

21

22

23

24

25

26

27

28

29

30

31

32

33

34

35

36

37

38

39

40

41

42

43

44

45

46

47

48

49

50

51

52

53

54

55

56

57

package com.itheima.rabbitmq.work;

import com.itheima.rabbitmq.util.ConnectionUtil;

import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer2 {

     public static void main(String[] args) throws Exception {

         Connection connection = ConnectionUtil.getConnection();

         // 创建频道

         Channel channel = connection.createChannel();

         // 声明(创建)队列

         /**

          * 参数1:队列名称

          * 参数2:是否定义持久化队列

          * 参数3:是否独占本次连接

          * 参数4:是否在不使用的时候自动删除队列

          * 参数5:队列其它参数

         */

         channel.queueDeclare(Producer.QUEUE_NAME, true , false , false , null );

         //一次只能接收并处理一个消息

         channel.basicQos( 1 );

         //创建消费者;并设置消息处理

         DefaultConsumer consumer = new DefaultConsumer(channel){

             @Override

             /**

              * consumerTag 消息者标签,在channel.basicConsume时候可以指定

              * envelope 消息包的内容,可从中获取消息id,消息routingkey,交换机,消息和重传标志(收到消息失败后是否需要重新发送)

              * properties 属性信息

              * body 消息

             */

             public void handleDelivery(String consumerTag, Envelope envelope,

                     AMQP.BasicProperties properties, byte [] body) throws IOException {

                 try {

                     //路由key

                     System.out.println( "路由key为:" + envelope.getRoutingKey());

                     //交换机

                     System.out.println( "交换机为:" + envelope.getExchange());

                     //消息id

                     System.out.println( "消息id为:" + envelope.getDeliveryTag());

                     //收到的消息

                     System.out.println( "消费者2-接收到的消息为:" + new String(body, "utf-8" ));

                     Thread.sleep( 1000 );

                     //确认消息

                     channel.basicAck(envelope.getDeliveryTag(), false );

                 } catch (InterruptedException e) {

                     e.printStackTrace();

                 }

             }

         };

         //监听消息

         /**

          * 参数1:队列名称

          * 参数2:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复会删除消息,设置为false则需要手动确认

          * 参数3:消息接收到后回调

         */

         channel.basicConsume(Producer.QUEUE_NAME, false , consumer);

     }

}

三、测试

启动两个消费者,然后再启动生产者发送消息;到IDEA的两个消费者对应的控制台查看是否竞争性的接收到消息。

总结

在一个队列中如果有多个消费者,那么消费者之间对于同一个消息的关系是 竞争 的关系。

到此这篇关于如何在RabbitMQ中实现Work queues模式的文章就介绍到这了,希望对你有所帮助,更多相关RabbitMQ内容请搜索以前的文章或继续浏览下面的相关文章,希望大家以后多多支持!

原文链接:https://blog.csdn.net/Java_Caiyo/article/details/115689824

查看更多关于在RabbitMQ中实现Work queues工作队列模式的详细内容...

  阅读:21次