SpringBoot整合RabbitMQ的基础搭建

RabbitMQ的整体概括

RabbitMQ是对于AMQP(高级消息队列协议)的具体实现,是一个用于在分布式系统中存储转发消息的网络通信协议。

应用解耦

此模型表示,消息中间件broker是消息的容器,它从生产者接收消息,并根据路由规则把消息投递给消费者。需要注意的是,生产者不是直接将消息投递给broker,而是通过路由键(RoutingKey)投递给交换机(Exchange),交换机绑定具体的队列(Queue),而消费者消费的消息只与队列有关,消费将会有推拉两种模式。

图示如下:

AMQP协议模型

AMQP 协议层角色相关的概念:

  • 生产者(producer):产生消息的应用,能够传递消息到消息中间件的应用。
  • 消息中间件(brokers):消息传递的中间载体,即我们今天的主角 RabbitMQ。
  • 消费者(consumers):接收并处理消息的应用。从消息中间件获取消息并处理。
  • 连接(Connection):生产者或消费者和消息中间件之间需要建立起连接。AMQP 应用层协议使用的是能够提供可靠投递的 TCP 连接,AMQP 的连接通常是长连接,AMQP 使用认证机制并且提供 TLS(SSL)保护。当我们的生产者 或 消费者 不再需要连接到消息中间件的的时候,需要优雅的释放掉它们与消息中间件 TCP 连接,而不是直接将 TCP 连接关闭。
  • 信道(channel):通常情况下生产者 或 消费者 需要与 消息中间件之间建立多个连接。无论怎样,同时开启多个 TCP 连接都是不合适的,因为这样做会消耗掉过多的系统资源。AMQP 协议提供了信道(channel)这个概念来处理多连接,可以把通道理解成共享一个 TCP 连接的多个轻量化连接。一个特定通道上的通讯与其他通道上的通讯是完全隔离的,因此每个 AMQP 方法都需要携带一个通道号,这样客户端就可以指定此方法是为哪个信道准备的。

消息中间件相关的概念:

  • 虚拟主机(vHosts):虚拟主机概念,一个 Virtual Host 里面可以有若干个 Exchange 和 Queue,我们可以控制用户在 Virtual Host 的权限。
  • 用户(User):最直接了当的认证方式,谁可以使用当前的消息中间件。
  • 交换机(Exchange):交换机接收生产者发出的消息并且路由到由交换机类型和被称作绑定(bindings)的规则所决定的到队列中,交换机不存储消息。
  • 消息(message):生产者产生的和消费者处理的消息。
  • 路由键(routing key):路由关键字,交换机 exchange 的路由规则利用这个关键字进行消息投递到消息队列。(路由键长度不能超过 255 个字节)
  • 绑定(Binding):Binding 可以理解为交换机 Exchange 路由消息到消息队列的路由规则关系(即消息队列和交换机的绑定)。当交换机 Exchange 收到生产者传递的消息 Message 时会解析其 Routing Key,Exchange 根据 Routing Key 与交换机类型 Exchange Type 将 Message 路由到消息队列中去。

SpringBoot整合RabbitMQ

简单了解了rabbitmq的概念,我们来实践一下与springboot的整合。
首先引入jar包。


<groupId>org.springframework.boot</groupId>

    <artifactId>spring-boot-starter-amqp</artifactId>

    <version>1.7.9.RELEASE</version>

</dependency>

配置连接工厂

    /**
     * 创建连接工厂
     * @return
     */
    @Bean(name = "adapterConnectionFactory")
    @Order(value = 2)
    public ConnectionFactory adapterConnectionFactory() {
        //创建连接工厂
        CachingConnectionFactory connectionFactory = new CachingConnectionFactory();
        //设置集群方式
        connectionFactory.setAddresses(rabbitMQProperties.getAddress());
        //connectionFactory.setHost(rabbitMQProperties.getHost());设置单节点方式
        //设置端口
        connectionFactory.setPort(rabbitMQProperties.getPort());
        //设置用户名
        connectionFactory.setUsername(rabbitMQProperties.getUsername());
        //设置密码
        connectionFactory.setPassword(rabbitMQProperties.getPassword());
        //设置虚拟主机
        connectionFactory.setVirtualHost(rabbitMQProperties.getVirtualHost());
        //消息确认机制confirm-callback或return-callback,成功后confirm,失败后回调
        connectionFactory.setPublisherReturns(true);
        connectionFactory.setPublisherConfirms(true);
        return connectionFactory;
    }

消息发送的组件RabbitTemplate

    /**
     * 创建消息发送组件
     * @return
     */
    @Bean(name = "adapterRabbitTemplate")
    public RabbitTemplate rabbitTemplate() {
        RabbitTemplate rabbitTemplate = new RabbitTemplate(adapterConnectionFactory());
        //exchange根据路由键匹配不到对应的queue时将会调用basic.return将消息返还给生产者
        rabbitTemplate.setMandatory(true);
        rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
            //消息成功发送到broker
            @Override
            public void confirm(CorrelationData correlationData, boolean ack, String cause) {
                ApiLog.info("mq message send (ACK)status =", ack);
            }
        });
        rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
             //消息发送失败
            @Override
            public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
                log.info("消息丢失:exchange({}),route({}),replyCode({}),replyText({}),message:{}", exchange, routingKey, replyCode, replyText, message);
            }
        });
        return rabbitTemplate;
    }

监听器容器工厂SimpleRabbitListenerContainerFactory

     /**
     * 消费者监听
     *
     * @return
     */
    @Bean(name = "singleListenerContainer")
    public SimpleRabbitListenerContainerFactory listenerContainer() {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(adapterConnectionFactory());
        factory.setMessageConverter(new Jackson2JsonMessageConverter());
        //单台并发消费者数量
        factory.setConcurrentConsumers(10);
        //单台并发消费的最大消费者数量
        factory.setMaxConcurrentConsumers(30);
        //预取消费数量,unacked数量超过这个值broker将不会接收消息
        factory.setPrefetchCount(5);
        //有事务时处理的消息数
        factory.setTxSize(1);
        //消息确认机制
        factory.setAcknowledgeMode(AcknowledgeMode.AUTO);
        //构建retryConfig,用于在JavaConfig的模式下读取并发参数
        RabbitProperties.AmqpContainer config = rabbitMQProperties.getListener().getSimple();
        RabbitProperties.ListenerRetry retryConfig = config.getRetry();
        RetryInterceptorBuilder<?> builder = (retryConfig.isStateless()
                ? RetryInterceptorBuilder.stateless()
                : RetryInterceptorBuilder.stateful());
        //最大重试次数,消费者异常之后
        builder.maxAttempts(retryConfig.getMaxAttempts());
        builder.backOffOptions(retryConfig.getInitialInterval(),
                retryConfig.getMultiplier(), retryConfig.getMaxInterval());
        MessageRecoverer recoverer = (this.messageRecoverer != null
                ? this.messageRecoverer : new RejectAndDontRequeueRecoverer());
        builder.recoverer(recoverer);
        factory.setAdviceChain(builder.build());
        return factory;
    }

重试参数说明,其他参数略

spring.rabbitmq.listener.simple.retry.enabled=true
//消费者异常之后的最大重试次数,JavaConfig方式需显示构建retryConfig
spring.rabbitmq.listener.simple.retry.max-attempts=4
spring.rabbitmq.listener.simple.retry.initial-interval=2000
spring.rabbitmq.listener.simple.default-requeue-rejected=true

最佳实践

初始化

 @Autowired
    private ConnectionFactory connectionFactory;

    @Value("${mq.queue.callback_queue}")
    private String callbackQueueKey;


    @Value("${mq.exchange}")
    private String exchange;

    @PostConstruct
    public void init(){
        RabbitAdmin admin = new RabbitAdmin(connectionFactory);
        //声明exchange
        Exchange topicExchange = new TopicExchange(exchange);
        admin.declareExchange(topicExchange);
        //声明queue
        Queue callbackQueue = new Queue(callbackQueueKey, true);
        admin.declareQueue(callbackQueue);
        //Binding
        admin.declareBinding(BindingBuilder.bind(callbackQueue)
                .to(topicExchange)
                .with(callbackQueueKey).noargs());
    }

生产者

        rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter());
        rabbitTemplate.setExchange(exchange);
        rabbitTemplate.setRoutingKey(callbackQueueKey);
        rabbitTemplate.convertAndSend(callBackRequest, new MessagePostProcessor() {
            @Override
            public Message postProcessMessage(Message message) throws AmqpException {
                MessageProperties properties = message.getMessageProperties();
                properties.setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                properties.setHeader(AbstractJavaTypeMapper.DEFAULT_CONTENT_CLASSID_FIELD_NAME, Obj.class);
                return message;
            }
        });

消费者

@RabbitListener(queues = "queueName", containerFactory = "singleListenerContainer")
    public void consumeMessage(@Payload Object obj, Channel channel, Message message) {
    doSometing...
}

需要注意的是,消费者的异常没有捕获或抛出,或者catch块里出现异常,在消息确认机制是AUTO的前提下将会无限重试进入死循环,这个时候可以设置最大重试次数或手动进行ack来处理。

如果需要手动ack,需要实现ChannelAwareMessageListener

@Override
public void onMessage(Message message, Channel channel) throws Exception {
        channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
    }

以上只是项目使用中的简单阐述,除此还有exchange的4种模式,死信队列,消息堆积,可靠性投递机制等待补充。

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 216,372评论 6 498
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,368评论 3 392
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 162,415评论 0 353
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,157评论 1 292
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,171评论 6 388
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,125评论 1 297
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,028评论 3 417
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,887评论 0 274
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,310评论 1 310
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,533评论 2 332
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,690评论 1 348
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,411评论 5 343
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,004评论 3 325
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,659评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,812评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,693评论 2 368
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,577评论 2 353

推荐阅读更多精彩内容

  • http://liuxing.info/2017/06/30/Spring%20AMQP%E4%B8%AD%E6%...
    sherlock_6981阅读 15,906评论 2 11
  • Spring Cloud为开发人员提供了快速构建分布式系统中一些常见模式的工具(例如配置管理,服务发现,断路器,智...
    卡卡罗2017阅读 134,651评论 18 139
  • 这个指导提供一个AMQP 0-9-1协议的概述,它是RabbitMq支持的一个协议。 什么是AMQP 0-9-1?...
    浪_6e80阅读 716评论 0 1
  • 什么叫消息队列? 消息(Message)是指在应用间传送的数据。消息可以非常简单,比如只包含文本字符串,也可以更复...
    Agile_dev阅读 2,371评论 0 24
  • AMQP大致内容就是,将消息和队列绑定起来,规定让进入到交换机中的具有某个路由键的消息进入到指定队列中去。 Rab...
    StevenMD阅读 1,858评论 0 3