RabbitMQ使用详解

RabbitMQ是基于AMQP的一款消息管理系统。AMQP(Advanced Message Queuing Protocol),是一个提供消息服务的应用层标准高级消息队列协议,其中RabbitMQ就是基于这种协议的一种实现。算法

常见mq:spring

  • ActiveMQ:基于JMS
  • RabbitMQ:基于AMQP协议,erlang语言开发,稳定性好
  • RocketMQ:基于JMS,阿里巴巴产品,目前交由Apache基金会
  • Kafka:分布式消息系统,高吞吐量

1 消息模型

RabbitMq有5种经常使用的消息模型网络

1.1 基本消息模型

这是最简单的消息模型,以下图:
app

生产者将消息发送到队列,消费者从队列中获取消息,队列是存储消息的缓冲区。负载均衡

再演示代码以前,咱们先建立一个工程rabbitmq-demo,并编写一个工具类,用于提供与mq服务建立链接异步

public class ConnectionUtil {
    /**
     * 创建与RabbitMQ的链接
     * @return
     * @throws Exception
     */
    public static Connection getConnection() throws Exception {
        //定义链接工厂
        ConnectionFactory factory = new ConnectionFactory();
        //设置服务地址
        factory.setHost("192.168.18.130");
        //端口
        factory.setPort(5672);
        //设置帐号信息,用户名、密码、vhost
        factory.setUsername("admin");
        factory.setPassword("admin");
        // 经过工程获取链接
        Connection connection = factory.newConnection();
        return connection;
    }
}

生产者发送消息

接下来是生产者发送消息,其过程包括:1.与mq服务创建链接,2.创建通道,3.声明队列(有相同队列则不建立,没有则建立),4.发送消息,代码以下:分布式

public class Send {
    private static final String QUEUE_NAME = "basic_queue";
    public static void main(String[] args) throws Exception {
        //消息发送端与mq服务建立链接
        Connection connection = ConnectionUtil.getConnection();
        //创建通道
        Channel channel = connection.createChannel();
        //声明队列
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        String message = "hello world";
        channel.basicPublish("", QUEUE_NAME, null, message.getBytes());
        System.out.println("生产者已发送:" + message);
        channel.close();
        connection.close();
    }
}

消费者接受消息

消费者在接收消息的过程须要经历以下几个步骤: 1.与mqfuwu创建链接,2.创建通道,3.声明队列,4,接收消息,代码以下:ide

public class Consumer1 {
    private static final String QUEUE_NAME = "basic_queue";
    public static void main(String[] args) throws Exception {
        //消息消费者与mq服务创建链接
        Connection connection = ConnectionUtil.getConnection();
        //创建通道
        Channel channel = connection.createChannel();
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            // 获取消息,而且处理,这个方法相似事件监听,若是有消息的时候,会被自动调用
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
                                       byte[] body) throws IOException {
                // body 即消息体
                String msg = new String(body);
                System.out.println("消费者1接收到消息:" + msg);
            }
        };
        channel.basicConsume(QUEUE_NAME, true, consumer);
    }
}

消息的接收与消费使用都须要在一个匿名内部类DefaultConsumer中完成spring-boot

注意:队列须要提早声明,若是未声明就使用队列,则会报错。若是不清楚生产者和消费者谁先声明,为了保证不报错,生产者和消费者都声明队列,队列的建立会保证幂等性,也就是说生产者和消费者都声明同一个队列,则只会建立一个队列工具

1.2 Work Queues工做队列模型

在基本消息模型中,一个生产者对应一个消费者,而实际生产过程当中,每每消息生产会发送不少条消息,若是消费者只有一个的话效率就会很低,所以rabbitmq有另一种消息模型,这种模型下,一个生产发送消息到队列,容许有多个消费者接收消息,可是一条消息只会被一个消费者获取。

生产者发送消息

与基本消息模型基本一致,这里测试循环发布20条消息:

public class Send {
    private static final String QUEUE_NAME = "work_queue";
    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        // 循环发布任务
        for (int i = 1; i <= 20; i++) {
            // 消息内容
            String message = "task .. " + i;
            channel.basicPublish("", QUEUE_NAME, null, message.getBytes());
            System.out.println("生产者发送消息:" + message);
            Thread.sleep(500);
        }
        channel.close();
        connection.close();
    }
}

消费者1

public class Consumer1 {
    private static final String QUEUE_NAME = "work_queue";
    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                String msg = new String(body);
                System.out.println("消费者1接收到消息:" + msg);
                try {
                    Thread.sleep(50);//模拟消费耗时
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        };
        channel.basicConsume(QUEUE_NAME, true, consumer);
    }
}

消费者2

public class Consumer2 {
    private static final String QUEUE_NAME = "work_queue";
    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,byte[] body) throws IOException {
                String msg = new String(body);
                System.out.println("消费者2接收到消息:" + msg);
                try {
                    Thread.sleep(50);//模拟消费耗时
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                channel.basicAck(envelope.getDeliveryTag(), false);
            }
        };
        channel.basicConsume(QUEUE_NAME, true, consumer);
    }
}

此时有两个消费者监听同一个队列,当两个消费者都工做时,生成者发送消息,就会按照负载均衡算法分配给不一样消费者,以下图:

1.3 订阅模型

在以前的模型中,一条消息只能被一个消费者获取,而在订阅模式中,能够实现一条消息被多个消费者获取。在这种模型下,消息传递过程当中比以前多了一个exchange交换机,生产者不是直接发送消息到队列,而是先发送给交换机,经由交换机分配到不一样的队列,而每一个消费者都有本身的队列:

解读:

一、1个生产者,多个消费者

二、每个消费者都有本身的一个队列

三、生产者没有将消息直接发送到队列,而是发送到了交换机

四、每一个队列都要绑定到交换机

五、生产者发送的消息,通过交换机到达队列,实现一个消息被多个消费者获取的目的

X(exchange)交换机的类型有如下几种:

Fanout:广播,交换机将消息发送到全部与之绑定的队列中去

Direct:定向,交换机按照指定的Routing Key发送到匹配的队列中去

Topics:通配符,与Direct大体相同,不一样在于Routing Key能够根据通配符进行匹配

注意:在发布订阅模型中,生产者只负责发消息到交换机,至于消息该怎么发,以及发送到哪一个队列,生产者都不负责。通常由消费者建立队列,而且绑定到交换机

订阅模型之Fanout

在广播模式下,消息发送的流程以下:

  1. 能够有多个消费者,每一个消费者都有本身的队列
  2. 每一个队列都要与exchange绑定
  3. 生产者发送消息到exchange
  4. exchange将消息把消息发送到全部绑定的队列中去
  5. 消费者从各自的队列中获取消息
生产者发送消息
public class Send {
    private static final String EXCHANGE_NAME = "fanout_exchange";
    
    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        // 声明exchange,指定类型为fanout
        channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT);
        String message = "hello world";
        channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());
        System.out.println("生产者发送消息:" + message);
        channel.close();
        connection.close();
    }
}
消费者
public class Consumer1 {
    private static final String QUEUE_NAME = "fanout_queue_1";
    private static final String EXCHANGE_NAME = "fanout_exchange";

    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        //消费者声明本身的队列
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        // 声明exchange,指定类型为direct
        channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT);
        //消费者将队列与交换机进行绑定
        channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "");
        channel.basicConsume(QUEUE_NAME, true, new DefaultConsumer(channel){
            @Override
            public void handleDelivery(String consumerTag,
                                       Envelope envelope,
                                       AMQP.BasicProperties properties,
                                       byte[] body)
                    throws IOException
            {
                String msg = new String(body);
                System.out.println("消费者1获取到消息:" + msg);
            }
        });
    }
}

其余消费者只需修改QUEUE_NAME便可

注意:exchange与队列同样都须要提早声明,若是未声明就使用交换机,则会报错。若是不清楚生产者和消费者谁先声明,为了保证不报错,生产者和消费者都声明交换机,一样的,交换机的建立也会保证幂等性。

订阅模型之Direct

在fanout模型中,生产者发布消息,全部消费者均可以获取全部消息。在路由模式(Direct)中,能够实现不一样的消息被不一样的队列消费,在Direct模式下,交换机再也不将消息发送给全部绑定的队列,而是根据Routing Key将消息发送到指定的队列,队列在与交换机绑定时会设定一个Routing Key,而生产者发送的消息时也须要携带一个Routing Key。

如图所示,消费者C1的队列与交换机绑定时设置的Routing Key是“error”, 而C2的队列与交换机绑定时设置的Routing Key包括三个:“info”,“error”,“warning”,假如生产者发送一条消息到交换机,并设置消息的Routing Key为“info”,那么交换机只会将消息发送给C2的队列。

生产者发送消息
public class Send {
    private static final String EXCHANGE_NAME = "direct_exchange";

    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        // 声明exchange,指定类型为direct
        channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
        String message = "新增一个订单";
        //生产者发送消息时,设置消息的Routing Key:"insert"
        channel.basicPublish(EXCHANGE_NAME, "insert", null, message.getBytes());
        System.out.println("生产者发送消息:" + message);
        channel.close();
        connection.close();
    }
}
消费者1
public class Consumer1 {
    private static final String QUEUE_NAME = "direct_queue_1";
    private static final String EXCHANGE_NAME = "direct_exchange";

    public static void main(String[] args) throws Exception {
        Connection connection = ConnectionUtil.getConnection();
        Channel channel = connection.createChannel();
        //消费者声明本身的队列
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        //声明交换机
        channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
        //消费者将队列与交换机进行绑定,而且设置Routing Key:"insert"
        channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "insert");
        channel.basicConsume(QUEUE_NAME, true, new DefaultConsumer(channel){
            @Override
            public void handleDelivery(String consumerTag,
                                       Envelope envelope,
                                       AMQP.BasicProperties properties,
                                       byte[] body)
                    throws IOException
            {
                String msg = new String(body);
                System.out.println("消费者1获取到消息:" + msg);
            }
        });
    }
}

其余消费者须要修改队列名QUEUE_NAME和Routing Key,上述生成者发送的消息,消费者1是能够获取到的

发布订阅之Topics

Topic类型的Exchange与Direct相比,都是能够根据RoutingKey把消息路由到不一样的队列。只不过Topic类型Exchange可让队列在绑定Routing key 的时候使用通配符

Routingkey 通常都是有一个或多个单词组成,多个单词之间以”.”分割,例如: item.insert

通配符规则:

#:匹配一个或多个词

     *:匹配很少很多刚好1个词

举例:

audit.#:可以匹配audit.irs.corporate 或者 audit.irs

     audit.*:只能匹配audit.irs

Topics生产者代码与Direct大体相同,只不过子声明交换机时,将类型设为BuiltinExchangeType.TOPIC(topic),

消费者代码也与Direct大体相同,也是在声明交换机时设置类型为topic,代码再也不演示

Spring AMQP

Spring AMQP是对AMQP的一种封装,目的是可以让咱们更简便的使用消息队列,下面介绍一下Spring AMQP在Spring boot中的使用方法

依赖和配置

添加AMQP的启动器:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

在application.yml中添加RabbitMQ的地址:

spring:
  rabbitmq:
    host: 192.168.18.130
    username: admin
    password: admin

消费者

消费者须要定义一个类,类中定义监听队列的方法

@Component
public class Listener {

    @RabbitListener(
            bindings = @QueueBinding(
                    value = @Queue(value = "spring.test.queue", durable = "false"),
                    exchange = @Exchange(value = "spring.test.exchange", type = ExchangeTypes.DIRECT),
                    key = "insert"
            )
    )
    public void listen(String msg){
        System.out.println("消费者接受到消息:" + msg);
    }
}

注解:

@Component:保证监听类被spring扫描到

@RabbitListener:

@RabbitListener包含不少内容,在发布订阅模式中,咱们可使用其中的“QueueBinding[] bindings”,其中QueueBinding底层以下:

其中Queue表示队列,Exchange表示交换机,key表示Routing Key

@RabbitListener(
            bindings = @QueueBinding(
                    value = @Queue(value = "spring.test.queue", durable = "false"),
                    exchange = @Exchange(value = "spring.test.exchange", type = ExchangeTypes.DIRECT),
                    key = "insert"
            )
    )

@Queue会建立队列

@Exchange会建立交换机

@QueueBinding会绑定队列和交换机

生产者发送消息

能够经过注解引入AmqpTemplate:

@RunWith(SpringRunner.class)
@SpringBootTest
public class SpringAMQPTest {
    @Resource
    private AmqpTemplate template;

    @Test
    public void testSendMsg() throws InterruptedException {
        String message = "hello spring";
        template.convertAndSend("spring.test.exchange", "insert", message);
        System.out.println("生产者发送消息:" + message);
        Thread.sleep(10000);//等待10s,让测试方法延迟结束,防止消费者将来得及获取消息
    }
}

RabbitMQ如何防止消息丢失

1. 消息确认机制(ACK)

RabbitMQ有一个ACK机制,消费者在接收到消息后会向mq服务发送回执ACK,告知消息已被接收。这种ACK分为两种状况:

  • 自动ACK:消息一旦被接收,消费者会自动发送ACK
  • 手动ACK:消息接收后,不会自动发送ACK,而是须要手动发送ACK

若是消费者没有发送ACK,则消息会一直保留在队列中,等待下次接收。但这里存在一个问题,就是一旦消费者发送了ACK,若是消费者后面宕机,则消息会丢失。所以自动ACK不能保证消费者在接收到消息以后可以正常完成业务功能,所以须要在消息被充分利用以后,手动ACK确认

自动ACK,basicConsume方法中将autoAck参数设为true便可:

手动ack,在匿名内部类中,手动发送ACK:

固然,若是设置了手动ack,但又不手动发送ACK确认,消息会一直停留在队列中,可能形成消息的重复获取

2. 持久化

消息确认机制(ACK)可以保证消费者不丢失消息,但假如消费者在获取消息以前mq服务宕机,则消息也会丢失,所以要保证消息在服务端不丢失,则须要将消息进行持久化。队列、交换机、消息都要持久化。

队列持久化

exchange持久化

消息持久化

3. 生产者确认

生成者在发送消息过程当中也可能出现错误或者网络延迟灯故障,致使消息未成功发送到交换机或者队列,或重复发送消息,为了解决这个问题,rabbitmq中有多个解决办法:

事务:

用事务将消息发送代码包围起来:

Confirm模式:

以下所示,在发送代码前执行channel.confirmSelect(),若是消息未正常发送,就会进入if代码块,能够进行重发也能够对失败消息进行记录

异步confirm方法:

顾名思义,就是生产者发送消息后不用等待服务端回馈发送状态,能够继续执行后面的代码,对于失败消息重发进行异步处理:

Spring AMQP中添加配置:

生产者确认机制,确保消息正确发送,若是发送失败会有错误回执,从而触发重试

spring:
  rabbitmq:
    publisher-confirms: true
相关文章
相关标签/搜索