ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

RabbitMQ exchange交换机机制

RabbitMQ exchange交换机机制 目录RabbitMQ 概念exchange交换机机制什么是交换机bindingDirect Exchange交换机Topic Exchange交换机Fanout Exchange交换机Header Exchange交换机RabbitMQ 的 Hello DemoSpring Boot实现RabbitMQ 的 Hello DemoSpring XML实现RabbitMQ 在生产环境下运用和出现的问题Spring RabbitMQ 注解消息的 JSON 传输消息持久化断线重连ACKRabbitMQ 概念RabbitMQ 即一个消息队列主要是用来实现应用程序的异步和解耦同时也能起到消息缓冲消息分发的作用。RabbitMQ使用的是AMQP协议它是一种二进制协议。默认启动端口 5672。在 RabbitMQ 中如下图结构rabbitmq左侧 P 代表 生产者也就是往 RabbitMQ 发消息的程序。中间即是 RabbitMQ其中包括了 交换机 和 队列。右侧 C 代表 消费者也就是往 RabbitMQ 拿消息的程序。那么其中比较重要的概念有 4 个分别为虚拟主机交换机队列和绑定。虚拟主机一个虚拟主机持有一组交换机、队列和绑定。为什么需要多个虚拟主机呢很简单RabbitMQ当中用户只能在虚拟主机的粒度进行权限控制。因此如果需要禁止A组访问B组的交换机/队列/绑定必须为A和B分别创建一个虚拟主机。每一个RabbitMQ服务器都有一个默认的虚拟主机“/”。交换机Exchange 用于转发消息但是它不会做存储如果没有 Queue bind 到 Exchange 的话它会直接丢弃掉 Producer 发送过来的消息。这里有一个比较重要的概念路由键。消息到交换机的时候交互机会转发到对应的队列中那么究竟转发到哪个队列就要根据该路由键。绑定也就是交换机需要和队列相绑定这其中如上图所示是多对多的关系。exchange交换机机制什么是交换机rabbitmq的message model实际上消息不直接发送到queue中中间有一个exchange是做消息分发producer甚至不知道消息发送到那个队列中去。因此当exchange收到message时必须准确知道该如何分发。是append到一定规则的queue还是append到多个queue中还是被丢弃这些规则都是通过exchagne的4种type去定义的。The core idea in the messaging model in RabbitMQ is that the producer never sends any messages directly to a queue. Actually, quite often the producer doesnt even know if a message will be delivered to any queue at all.Instead, the producer can only send messages to an exchange. An exchange is a very simple thing. On one side it receives messages from producers and the other side it pushes them to queues. The exchange must know exactly what to do with a message it receives. Should it be appended to a particular queue? Should it be appended to many queues? Or should it get discarded. The rules for that are defined by the exchange type.exchange是一个消息的agent每一个虚拟的host中都有定义。它的职责是把message路由到不同的queue中。bindingexchange和queue通过routing-key关联这两者之间的关系是就是binding。如下图所示,X表示交换机红色表示队列交换机通过一个routing-key去binding一个queuerouting-key有什么作用呢看Direct exchange类型交换机。Directed Exchange路由键exchange该交换机收到消息后会把消息发送到指定routing-key的queue中。那消息交换机是怎么知道的呢其实producer deliver消息的时候会把routing-key add到 message header中。routing-key只是一个messgae的attribute。A direct exchange delivers messages to queues based on a message routing key. The routing key is a message attribute added into the message header by the producer. The routing key can be seen as an address that the exchange use to decide how to route the message. A message goes to the queue(s) whose binding key exactly matches the routing key of the message.Default Exchange这种是特殊的Direct Exchange是rabbitmq内部默认的一个交换机。该交换机的name是空字符串所有queue都默认binding 到该交换机上。所有binding到该交换机上的queuerouting-key都和queue的name一样。Topic Exchange通配符交换机exchange会把消息发送到一个或者多个满足通配符规则的routing-key的queue。其中*表号匹配一个word#匹配多个word和路径路径之间通过.隔开。如满足a.*.c的routing-key有a.hello.c满足#.hello的routing-key有a.b.c.helo。Fanout Exchange扇形交换机该交换机会把消息发送到所有binding到该交换机上的queue。这种是publisher/subcribe模式。用来做广播最好。所有该exchagne上指定的routing-key都会被ignore掉。The fanout copies and routes a received message to all queues that are bound to it regardless of routing keys or pattern matching as with direct and topic exchanges. Keys provided will simply be ignored.Header Exchange设置header attribute参数类型的交换机。RabbitMQ 的 Hello DemoSpring Boot实现安装就不说了建议按照官方文档上做。先贴代码稍后解释代码如下配置交换机、队列、绑定和消息监听容器Configuration Data public class RabbitMQConfig { final static String queueName spring-boot; Bean Queue queue() { return new Queue(queueName, false); } Bean TopicExchange exchange() { return new TopicExchange(spring-boot-exchange); } Bean Binding binding(Queue queue, TopicExchange exchange) { return BindingBuilder.bind(queue).to(exchange).with(queueName); } Bean SimpleMessageListenerContainer container(ConnectionFactory connectionFactory, MessageListenerAdapter listenerAdapter) { SimpleMessageListenerContainer container new SimpleMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setQueueNames(queueName); container.setMessageListener(listenerAdapter); return container; } Bean Receiver receiver() { return new Receiver(); } Bean MessageListenerAdapter listenerAdapter(Receiver receiver) { return new MessageListenerAdapter(receiver, receiveMessage); } }配置接收者消费者public class Receiver { private CountDownLatch latch new CountDownLatch(1); public void receiveMessage(String message) { System.out.println(Received message ); latch.countDown(); } public CountDownLatch getLatch() { return latch; } }配置发送者生产者RestController public class Test { Autowired RabbitTemplate rabbitTemplate; RequestMapping(value /test/{abc},method RequestMethod.GET) public String test(PathVariable(value abc) String abc){ rabbitTemplate.convertAndSend(spring-boot, abc from RabbitMQ!); return abc; } }以上便可实现一个简单的 RabbitMQ Demo具体代码在点这里代码分析那么这里分为三个部分分析发消息交换机队列收消息。发送消息我们一般可以使用 RabbitTemplate这个是 Spring 封装给了我们便于我们发送信息我们调用rabbitTemplate.convertAndSend(spring-boot, xxx);即可发送信息。交换机队列如上代码我们需要配置交换机TopicExchange配置队列Queue并且配置他们之间的绑定Binding接收消息首先需要创建一个消息监听容器然后把我们的接受者注册到该容器中这样队列中有信息那么就会调用接收者的对应的方法。如上代码container.setMessageListener(listenerAdapter);其中MessageListenerAdapter 可以看做是 我们接收者的一个包装类new MessageListenerAdapter(receiver, receiveMessage);指明了如果有消息来那么调用接收者哪个方法进行处理。RabbitMQ 的 Hello DemoSpring XML实现spring xml方式实现RabbitMQ简单可读性较好配置简单配置和实现如下所示。配置文件上文已经讲述了rabbitmq的配置xml方式通过properites文件存放用户配置信息mq.host127.0.0.1 mq.usernameguest mq.passwordguest mq.port5672XML配置配置application-mq.xml配置文件声明连接、交换机、queue以及consumer监听。?xml version1.0 encodingUTF-8? beans xmlnshttp://www.springframework.org/schema/beans xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xmlns:rabbithttp://www.springframework.org/schema/rabbit xmlns:contexthttp://www.springframework.org/schema/context xsi:schemaLocationhttp://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-3.0.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context-3.0.xsd http://www.springframework.org/schema/rabbit http://www.springframework.org/schema/rabbit/spring-rabbit-1.0.xsd descriptionrabbitmq 连接服务配置/description !-- 连接配置 -- context:property-placeholder locationclasspath:mq.properties / rabbit:connection-factory idconnectionFactory host${mq.host} username${mq.username} password${mq.password} port${mq.port}/ rabbit:admin connection-factoryconnectionFactory/ !-- spring template声明-- rabbit:template exchangeamqpExchange idamqpTemplate connection-factoryconnectionFactory / !--申明queue-- rabbit:queue idtest_queue_key nametest_queue_key durabletrue auto-deletefalse exclusivefalse / !--申明exchange交换机并绑定queue-- rabbit:direct-exchange nameamqpExchange durabletrue auto-deletefalse idamqpExchange rabbit:bindings rabbit:binding queuetest_queue_key keytest_queue_key/ /rabbit:bindings /rabbit:direct-exchange !--consumer配置监听-- bean idreveiver classcom.demo.mq.receive.Reveiver / rabbit:listener-container connection-factoryconnectionFactory acknowledgeauto rabbit:listener queuestest_queue_key refreveiver methodreceiveMessage/ /rabbit:listener-container /beans上述代码中引入properties文件就不多说了。rabbit:connection-factory标签声明创建connection的factory工厂。rabbit-template声明spring template和上文spring中使用template一样。template可声明exchange。rabbit:queue声明一个queue并设置queue的配置项直接看标签属性就可以明白queue的配置项。rabbit:direct-exchange声明交换机并绑定queue。rabbit:listener-container申明监听container并配置consumer和监听routing-key。主配置文件剩下就简单了application-context.xml中把rabbitmq配置import进去。?xml version1.0 encodingUTF-8? beans xmlnshttp://www.springframework.org/schema/beans xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xmlns:taskhttp://www.springframework.org/schema/task xmlns:contexthttp://www.springframework.org/schema/context xmlns:aophttp://www.springframework.org/schema/aop xsi:schemaLocationhttp://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context-3.0.xsd http://www.springframework.org/schema/task http://www.springframework.org/schema/task/spring-task-3.0.xsd http://www.springframework.org/schema/aop http://www.springframework.org/schema/aop/spring-aop-3.0.xsd context:component-scan base-packagecom.demo.** / import resourceapplication-mq.xml / /beans生产者实现Producer实现发送消息还是使用template的convertAndSend() deliver消息。Service public class Producer { Autowired private AmqpTemplate amqpTemplate; private final static Logger logger LoggerFactory.getLogger(Producer.class); public void sendDataToQueue(String queueKey, Object object) { try { amqpTemplate.convertAndSend(queueKey, object); } catch (Exception e) { e.printStackTrace(); logger.error(exeception{},e); } } }消费者实现配置consumerpackage com.demo.mq.receive; import org.springframework.stereotype.Service; import java.util.concurrent.CountDownLatch; Service public class Reveiver { private CountDownLatch latch new CountDownLatch(1); public void receiveMessage(String message) { System.out.println(reveice msg message.toString()); latch.countDown(); } }测试deliver消息Controller RequestMapping(/demo/) public class TestController { private final static Logger logger LoggerFactory.getLogger(TestController.class); Resource private Producer producer; RequestMapping(/test/{msg}) public String send(PathVariable(msg) String msg){ logger.info(#TestController.send#abc{msg}, msg); System.out.println(msgmsg); producer.sendDataToQueue(test_queue_key,msg); return index; } }RabbitMQ 在生产环境下运用和出现的问题在生产环境中由于 Spring 对 RabbitMQ 提供了一些方便的注解所以首先可以使用这些注解。例如EnableRabbitEnableRabbit 和 Configuration 注解在一个类中结合使用如果该类能够返回一个 RabbitListenerContainerFactory 类型的 bean那么就相当于能够把该终端消费端和 RabbitMQ 进行连接。Ps生成端不是通过 RabbitListenerContainerFactory 来和 RabbitMQ 连接而是通过 RabbitTemplate RabbitListener当对应的队列中有消息的时候该注解修饰下的方法会被执行。RabbitHandler接收者可以监听多个队列不同的队列消息的类型可能不同该注解可以使得不同的消息让不同方法来响应。具体这些注解的使用可以参考这里的代码点这里首先生产环境下的 RabbitMQ 可能不会在生产者或者消费者本机上所以需要重新定义 ConnectionFactory即Bean ConnectionFactory connectionFactory() { CachingConnectionFactory connectionFactory new CachingConnectionFactory(host, port); connectionFactory.setUsername(userName); connectionFactory.setPassword(password); connectionFactory.setVirtualHost(vhost); return connectionFactory; }在该代码中new RabbitTemplate(connectionFactory);设置了生产端连接到RabbitMQtemplate.setMessageConverter(integrationEventMessageConverter());设置了 生产端发送给交换机的消息是以什么格式的在integrationEventMessageConverter()代码中public MessageConverter integrationEventMessageConverter() { Jackson2JsonMessageConverter messageConverter new Jackson2JsonMessageConverter(); return messageConverter; }如上Jackson2JsonMessageConverter指明了 JSON。上述代码的最后template.setExchange(exchangeName);指明了 要把生产者要把消息发送到哪个交换机上。有了上述那么我们即可使用rabbitTemplate.convertAndSend(spring-boot, xxx);发送消息xxx 表示任意类型因为上述的设置会帮我们把这些类型转化成 JSON 传输。接着生产端发送我们说过了那么现在可以看看消费端对于消费端我们可以只创建SimpleRabbitListenerContainerFactory它能够帮我们生成 RabbitListenerContainer然后我们再使用 RabbitListener 指定接收者收到信息时处理的方法。Bean(namemyListenContainer) public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory() { SimpleRabbitListenerContainerFactory factory new SimpleRabbitListenerContainerFactory(); factory.setMessageConverter(integrationEventMessageConverter()); factory.setConnectionFactory(connectionFactory()); return factory; }这其中factory.setMessageConverter(integrationEventMessageConverter());指定了我们接受消息的时候以 JSON 传输的消息可以转换成对应的类型传入到方法中。例如Slf4j Component RabbitListener(containerFactory helloRabbitListenerContainer,queues spring-boot) public class Receiver { RabbitHandler public void receiveTeacher(Teacher teacher) { log.info(##### {},teacher); } }可能出现的问题消息持久化在生产环境中我们需要考虑万一生产者挂了消费者挂了或者 rabbitmq 挂了怎么样。一般来说如果生产者挂了或者消费者挂了其实是没有影响因为消息就在队列里面。那么万一 rabbitmq 挂了之前在队列里面的消息怎么办其实可以做消息持久化RabbitMQ 会把信息保存在磁盘上。做法是可以先从 Connection 对象中拿到一个 Channel 信道对象然后再可以通过该对象设置 消息持久化。生产者或者消费者断线重连这里 Spring 有自动重连机制。ACK 确认机制每个Consumer可能需要一段时间才能处理完收到的数据。如果在这个过程中Consumer出错了异常退出了而数据还没有处理完成那么 非常不幸这段数据就丢失了。因为我们采用no-ack的方式进行确认也就是说每次Consumer接到数据后而不管是否处理完 成RabbitMQ Server会立即把这个Message标记为完成然后从queue中删除了。如果一个Consumer异常退出了它处理的数据能够被另外的Consumer处理这样数据在这种情况下就不会丢失了注意是这种情况下。为了保证数据不被丢失RabbitMQ支持消息确认机制即acknowledgments。为了保证数据能被正确处理而不仅仅是被Consumer收到那么我们不能采用no-ack。而应该是在处理完数据后发送ack。在处理数据后发送的ack就是告诉RabbitMQ数据已经被接收处理完成RabbitMQ可以去安全的删除它了。如果Consumer退出了但是没有发送ack那么RabbitMQ就会把这个Message发送到下一个Consumer。这样就保证了在Consumer异常退出的情况下数据也不会丢失。总结RabbitMQ 作用异步解耦缓冲消息分发。RabbitMQ 主要分为3个部分生产者交换机和队列消费者。需要注意消息持久化目的为了防止 RabbitMQ 宕机考虑 ACK 机制目的为了如果消费者对消息的处理失败了那么后续要如何处理。写在最后写出来说出来才知道对不对知道不对才能改正改正了才能成长。在技术方面希望大家眼里都容不得沙子。如果有不对的地方或者需要改进的地方希望可以指出万分感谢。
返回列表