美文网首页
分析开源代码之消息队列

分析开源代码之消息队列

作者: 只不过33 | 来源:发表于2020-05-25 15:02 被阅读0次

    分析一下别人写的开源代码吧

    别人写的开源代码的地址:IOTGate:https://gitee.com/pnoker/iot-dc3

    简介一下其中一个子模块的文件划分吧。

    api:存放与外界交互的程序文件。

    config:存放各“保姆”的配置程序文件,这些配置程序文件的作用是向spring注入各“保姆”的所需的一些对象。

    service:存放的程序文件的作用是为api中的程序提供代码接口 

    别人写的开源代码的地址:https://gitee.com/yidao620/springboot-bucket?_from=gitee_search

    package com.xncoding.pos.config;

    import org.slf4j.Logger;

    import org.slf4j.LoggerFactory;

    import org.springframework.amqp.core.*;

    import org.springframework.amqp.rabbit.core.RabbitTemplate;

    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;

    import org.springframework.beans.factory.annotation.Qualifier;

    import org.springframework.context.annotation.Bean;

    import org.springframework.context.annotation.Configuration;

    import javax.annotation.Resource;

    /**

    * RabbitConfig

    *

    * @author XiongNeng

    * @version 1.0

    * @since 2018/3/1

    */

    @Configuration

    public class RabbitConfig {

        @Resource

        private RabbitTemplate rabbitTemplate;

        /**

        * 定制化amqp模版      可根据需要定制多个

        * <p>

        * <p>

        * 此处为模版类定义 Jackson消息转换器

        * ConfirmCallback接口用于实现消息发送到RabbitMQ交换器后接收ack回调  即消息发送到exchange  ack

        * ReturnCallback接口用于实现消息发送到RabbitMQ 交换器,但无相应队列与交换器绑定时的回调  即消息发送不到任何一个队列中  ack

        *

        * @return the amqp template

        */

        // @Primary

        @Bean

        public AmqpTemplate amqpTemplate() {

            Logger log = LoggerFactory.getLogger(RabbitTemplate.class);

            // 使用jackson 消息转换器

            rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter());

            rabbitTemplate.setEncoding("UTF-8");

            // 消息发送失败返回到队列中,yml需要配置 publisher-returns: true

            rabbitTemplate.setMandatory(true);

            rabbitTemplate.setReturnCallback((message, replyCode, replyText, exchange, routingKey) -> {

                String correlationId = message.getMessageProperties().getCorrelationIdString();

                log.debug("消息:{} 发送失败, 应答码:{} 原因:{} 交换机: {}  路由键: {}", correlationId, replyCode, replyText, exchange, routingKey);

            });

            // 消息确认,yml需要配置 publisher-confirms: true

            rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {

                if (ack) {

                    log.debug("消息发送到exchange成功,id: {}", correlationData.getId());

                } else {

                    log.debug("消息发送到exchange失败,原因: {}", cause);

                }

            });

            return rabbitTemplate;

        }

        /* ----------------------------------------------------------------------------Direct exchange test--------------------------------------------------------------------------- */

        /**

        * 声明Direct交换机 支持持久化.

        *

        * @return the exchange

        */

        @Bean("directExchange")

        public Exchange directExchange() {

            return ExchangeBuilder.directExchange("DIRECT_EXCHANGE").durable(true).build();

        }

        /**

        * 声明一个队列 支持持久化.

        *

        * @return the queue

        */

        @Bean("directQueue")

        public Queue directQueue() {

            return QueueBuilder.durable("DIRECT_QUEUE").build();

        }

        /**

        * 通过绑定键 将指定队列绑定到一个指定的交换机 .

        *

        * @param queue    the queue

        * @param exchange the exchange

        * @return the binding

        */

        @Bean

        public Binding directBinding(@Qualifier("directQueue") Queue queue,

                                    @Qualifier("directExchange") Exchange exchange) {

            return BindingBuilder.bind(queue).to(exchange).with("DIRECT_ROUTING_KEY").noargs();

        }

        /* ----------------------------------------------------------------------------Fanout exchange test--------------------------------------------------------------------------- */

        /**

        * 声明 fanout 交换机.

        *

        * @return the exchange

        */

        @Bean("fanoutExchange")

        public FanoutExchange fanoutExchange() {

            return (FanoutExchange) ExchangeBuilder.fanoutExchange("FANOUT_EXCHANGE").durable(true).build();

        }

        /**

        * Fanout queue A.

        *

        * @return the queue

        */

        @Bean("fanoutQueueA")

        public Queue fanoutQueueA() {

            return QueueBuilder.durable("FANOUT_QUEUE_A").build();

        }

        /**

        * Fanout queue B .

        *

        * @return the queue

        */

        @Bean("fanoutQueueB")

        public Queue fanoutQueueB() {

            return QueueBuilder.durable("FANOUT_QUEUE_B").build();

        }

        /**

        * 绑定队列A 到Fanout 交换机.

        *

        * @param queue          the queue

        * @param fanoutExchange the fanout exchange

        * @return the binding

        */

        @Bean

        public Binding bindingA(@Qualifier("fanoutQueueA") Queue queue,

                                @Qualifier("fanoutExchange") FanoutExchange fanoutExchange) {

            return BindingBuilder.bind(queue).to(fanoutExchange);

        }

        /**

        * 绑定队列B 到Fanout 交换机.

        *

        * @param queue          the queue

        * @param fanoutExchange the fanout exchange

        * @return the binding

        */

        @Bean

        public Binding bindingB(@Qualifier("fanoutQueueB") Queue queue,

                                @Qualifier("fanoutExchange") FanoutExchange fanoutExchange) {

            return BindingBuilder.bind(queue).to(fanoutExchange);

        }

    }

    amqp这个“保姆”是干嘛的呢?

    跟管子又有什么关系呢?当管子运输客户端说的话的时候,可能客户端说话又多又快,一时忙不过来也是有可能的,但是话不会等待啊?

    那该怎么办呢?虽然话不能等待,但是我们可以找个“保姆”来负责这事儿,让话能够“等待”。这样的保姆存在吗?存在的,她就是 Rabbitpq。嗯,嗯就是这样的一个“保姆”做到了这样的工作。

    可以预想,当一句话过来后,可能先存到这个"保姆"提供的内存中。

    然后“保姆”,会有秩序地再将这些话发给服务器。

    事实上这个“保姆”管理的管子多了去了,连服务器和数据库的管子也会管,当然服务器和服务器之间的管子也会管,一个“模块”和另一个“模块”之间的管子也会管。

    这些管子已经不是原来那种单一的管子了,而是各种各样的管子。

    package com.xncoding.pos.mq;

    import com.rabbitmq.client.Channel;

    import org.slf4j.Logger;

    import org.slf4j.LoggerFactory;

    import org.springframework.amqp.core.Message;

    import org.springframework.amqp.rabbit.annotation.RabbitListener;

    import org.springframework.stereotype.Component;

    import java.io.IOException;

    /**

    * 消息监听器

    *

    * @author XiongNeng

    * @version 1.0

    * @since 2018/3/1

    */

    @Component

    public class Receiver {

        private static final Logger log = LoggerFactory.getLogger(Receiver.class);

        /**

        * FANOUT广播队列监听一.

        *

        * @param message the message

        * @param channel the channel

        * @throws IOException the io exception  这里异常需要处理

        */

        @RabbitListener(queues = {"FANOUT_QUEUE_A"})

        public void on(Message message, Channel channel) throws IOException {

            channel.basicAck(message.getMessageProperties().getDeliveryTag(), true);

            log.debug("FANOUT_QUEUE_A " + new String(message.getBody()));

        }

        /**

        * FANOUT广播队列监听二.

        *

        * @param message the message

        * @param channel the channel

        * @throws IOException the io exception  这里异常需要处理

        */

        @RabbitListener(queues = {"FANOUT_QUEUE_B"})

        public void t(Message message, Channel channel) throws IOException {

            channel.basicAck(message.getMessageProperties().getDeliveryTag(), true);

            log.debug("FANOUT_QUEUE_B " + new String(message.getBody()));

        }

        /**

        * DIRECT模式.

        *

        * @param message the message

        * @param channel the channel

        * @throws IOException the io exception  这里异常需要处理

        */

        @RabbitListener(queues = {"DIRECT_QUEUE"})

        public void message(Message message, Channel channel) throws IOException {

            channel.basicAck(message.getMessageProperties().getDeliveryTag(), true);

            log.debug("DIRECT " + new String(message.getBody()));

        }

    }

    消息队列 本质上应该是一种转发  当然转发的时候会夹带“私货” :给这些访问排个队

    但是 转发的时候有针对性的进行转发 

    相关文章

      网友评论

          本文标题:分析开源代码之消息队列

          本文链接:https://www.haomeiwen.com/subject/wcfsahtx.html