显示标签为“消息中间件”的博文。显示所有博文
显示标签为“消息中间件”的博文。显示所有博文

2013年6月18日星期二

通过消息中间件和云计算实现系统可伸缩性


文章可以转载, 但需要以超链接形式标明文章原始出处 http://blogs.huihoo.com

Messaging, 消息传递, 消息队列, 即时通信, 统一消息正获得越来越多的关注, 消息处理在互联网和企业生产系统中扮演着极其重要的角色.
Twitter 的成功也将实时消息处理推到了一个新高度. Google Wave 也正努力把邮件和消息完美地整合在一起. 越来越多的互联网企业都渴望打造一个类似Twitter的消息基础设施.
Twitter的消息队列Kestrel使 用Scala编 写.  惊讶的是其核心就只有1500行Scala代码.
我们再看看有没有其他的一些做法, 在这里我们想更多考虑互联网和企业内部的不同应用场景以及它们之间未来的融合.
XMPP 和 AMQP 是两个开放的消息标准:
  • Extensible Messaging and Presence Protocol (XMPP) 是基于可扩展标记语言(XML)的协议, 它用于即时消息(IM)以及状态显示(Presence)
  • Advanced Message Queuing Protocol(AMQP) is an Open Standard for Messaging Middleware.
ejabberd 和 RabbitMQ 是两个开源的消息中间件, 它们也有了越来越多的成功应用.
  • ejabberd is a Jabber/XMPP instant messaging server.
  • RabbitMQ is an implementation of AMQP.
ejabberd  和 RabbitMQ 都使用 Erlang 语言开发, 通过网关它们能很好得集成在一起. 为企业打造一个强大的消息基础设施. 我们建议将它们部署在”云”中, 如 Amazon EC2. 以获得足够的伸缩性.


通过 RabbitHub 能连通更多的开放标准, 形成一个更加开放、统一的消息平台.

又拍网架构 -- 前端PHP后台Python +消息中间件 RabbitMQ + 分库步骤

又拍网的服务器端开发语言主要是PHP和Python,其中PHP用于编写Web逻辑(通过HTTP和用户直接打交道), 而Python则主要用于开发内部服务和后台任务。 这是目前Web2.0公司的通用选项。

消息中间件:RabbitMQ
1.RabbitMQ: high performance messaging solution
3.通过消息中间件和云计算实现系统可伸缩性
5.RabbitMQ安装和测试小记


分库的一般阶段:
1. 一台主库和一台从库组成。 从库只用作备份和容灾,当主库出现故障时,从库就手动变成主库,一般情况下,从库不作读写操作(同步除外)。
2. 一台主库和多台从库。一些实时性要求不高的Query放到从库去执行。后面又通过添加多个从库来分流查询压力。
3. 随着数据量的增加,主库的写压力也越来越大,这是要拆库:
   3.1. 垂直拆分:是指按功能模块拆分,比如可以将群组相关表和照片相关表存放在不同的数据库中,这种方式多个数据库之间的表结构不同。
   3.2. 水平拆分:而水平拆分是将同一个表的数据进行分块保存到不同的数据库中,这些数据库中的表结构完全相同。
   3.3. 一般都会先进行垂直拆分,因为这种方式拆分方式实现起来比较简单,根据表名访问不同的数据库就可以了。

2013年6月17日星期一

基于消息的分布式架构设计

背景:

随着社会的发展,经济的飞跃,传统的单系统模式(webApp+DB)已经很难满足业务场景的需要。企业系统开始不断演化成多个子系统并存协作的局面。大大降低了系统间的耦合性,更重要的便于子系统的扩展、升级、维护等。

谈到系统间的协作,目前常用两种方式:
1、基于Http协议
通过客户端发起的get、post请求,服务端接收request请求,处理请求,得到响应内容,通过网络传送到客户端,由浏览器解析出一个可视化的页面。
这种交互最大的优势是实时性,通过HTTP请求连接各个子系统,从而跨服务器来完成一个完整的业务流程。缺点协议请求头的信息较少,一般都是关键参数,完整数据由下一个子系统从数据库、文件系统来获取,从来保证前后的业务数据衔接。

2、基于消息的模式。
这种模式一个很重要前提是对实时性要求不高。优点可以有效降低模块的耦合性,减轻主干业务流程,将大量的业务交由后台任务来处理,有效缩短系统响应时间,提高系统TPS。
比如用户下单成功后发送邮件功能,属于非主干功能,完全可以从下单的主干业务逻辑剥离出来,从来提高下单的响应速度。而发送邮件的功能则由邮件服务器接收异步消息来跟踪处理,带有点分布式集群的感觉,
将一个任务有效拆分到多台服务器来完成。


所谓消息本质上是一种数据结构(当然,对象也可以看做是一种特殊的消息),它包含生产者与消费者双方都能识别的数据,这些数据需要在不同的服务器之间进行传递,并可能会被多个完全不同的客户端消费。



消息队列降低了生产者和消费者之间的耦合性,他们不会存在直接的代码依赖,方便各自的扩展,比如生产者因为业务下线,导致代码下线,而消费端不用同时跟进处理,只是队列不会有消息,这样方便于更加灵活的协调开发资源,而不必一方下线,所有的依赖全部受影响,产生较高维护成本。另外我们也可以随意对生产者和消费者扩展,引入多个消息队列,他们之间的依赖可以配置在XML文件中,通过JNDI来获取消息队列Queue,每次加载时,通过lookup服务首先通过读取配置文件来获取通道。

常见的消息模型分为:点对点模型;发布-订阅模型
点对点模型:Point to Point,消息被生产者放到一个队列中,消费者从消息队列中取走消息。消息一旦被一个消费者取走后,消息就从队列中移除。这意味着即使有多个消费观察一个队列,但一个消息只能被一个消费者取走。
发布-订阅模型:Publish/Subscribe,发布者发布一条消息可以发送给所有的订阅用户,所有的订阅用户都有处理某一条消息的机会。
对于订阅者而言,有两种处理消息的方式。一种是广播机制,这时消息通道中的消息在出列的同时,还需要复制消息对象,将消息传递给多个订阅者。例如,有多个子系统都需要获取从CRM系统传来的客户信息,并根据传递过来的客户信息,进行相应的处理。此时的消息通道又被称为Propagation通道。另一种方式则属于抢占机制,它遵循同步方式,在同一时间只能有一个订阅者能够处理该消息。实现Publisher-Subscriber模式的消息通道会选择当前空闲的唯一订阅者,并将消息出列,并传递给订阅者的消息处理方法。
目前使用较多的是广播机制的消息处理方式,且将topic与queue有效组合
一个生产消息的事件对应一个topic,topic下面可以挂多个queue,当然一个queue也可以挂在多个topic下面,每个queue都对应一个消息的消费端,唯一消费,保证消费的准确性。
如下图所示,当下单时,会将下单的相关信息封装到消息体中,发送到下单事件关联的那个topic1中,然后Topic会将消息复制发送到挂载在其下面的所有队列上,将Message复制到快照队列、成交记录统计队列中,消息端会监听队列,
如果有消息 ,则启动任务线程,来进行相关的业务处理。
在引入消息队列时重点要注意以下几点:
  • 并发:选择的消息队列一定要很好地支持用户访问的并发性;
  • 安全:消息队列是否提供了足够的安全机制;
  • 性能伸缩:不能让消息队列成为整个系统的单一性能瓶颈;
  • 部署:尽可能让消息队列的部署更为容易;
  • 灾备:不能因为意外的错误、故障或其他因素导致处理数据的丢失,最好可以写入磁盘,持久化存储;
  • API易用性:处理消息的API必须足够简单、并能够很好地支持测试与扩展
  • 容量:队列的容量一定要大,至少可以存储千万级别的消息体

目前市场上有很多成熟的消息框架:如Active MQ,IBM 的MQ,JBoss MQ,MSMQ等,各有各的优势,在使用前一定要充分衡量是否可以满足自己的业务需求
下图是napoli的client类图:

JMS与Java消息中间件


 JMS是一种企业消息传送的API,并不是MOM消息中间件系统的全部,JMS也是一种规范,类似于JDBC,我们通JMS API访问 JMS的服务器。目前JMS服务器的主要产品有 IBM WebSphere MQ、SonicMQ、Sun Open MQ、BEA  WebLogic  JMS、Oracle AQ  以及 我们最常用的JBoss MQ 和 Apache的 ActiveMQ.。而JMS服务器是MQ(Message Queen)产品家族中的一种,Microsoft Message Queuing(mSMQ)也是MQ产品类似于JMS服务器。
JMS在J2EE系统中的应用场景
   异步消息比同步消息操作更加便利。在J2EE系统中最常见的一个场景,前端浏览器 jsp/servlet向服务器端发出一个请求,jsp/servlet将请求传递给后端的应用程序处理业务逻辑,业务模块将响应的结果直接返回给客户端,而不是真正的计算结果,例如一个网站的用户注册功能,一个用户点击注册以后,将会发送一份邮件给他当时的注册邮箱,如果需要等到邮件发送成功再返回给用户结果的话,用户体验将会很差,所以将结果直接返回给用户,将用户注册的信息通过消息发送给后端程序慢慢处理。
JMS Cluster
谈到 JMS一词 大体上有 3个部分  1 消息发送端  2中间件服务器  3消息接收端  3个组件缺一不可。 JMS消息分为两种消息模式,点对点和发布者/订阅者。许多提供商支持这一因此,程序员可以在他们的分布式软件中实现面向消息的操作,这些操作将具有不同面向消息中间件产品的可移植性。
 

http://l99raa.bay.livefilestore.com/y1pg92LAov90l3hJaPBGkJnXsqRwNHnmdw7vpbm6VozavIPtJJDnQzKqxtpPwrB6DjlEitcPYSr5ON62WkX5KmDzreapKUaMet6/DRF.png?psid=1
Java消息服务器是指,将数据通过消息作为载体在网络中从一个系统异步传送给另一个系统。这样的异步消息传送意味着:发送者不需要等待接收者接收或处理该消息;它可以自由地发送消息并持续进行处理。这样一个异步式的架构主要依赖于一台消息服务器(message server)。消息服务器,也称为消息路由器(message router)或代理(broker),它负责从一个消息传送客户端向其他消息传送客户端传送消息。
JMS对与一个大型系统是必不可少的一个应用组件,可以利用JMS来实现3个目的:
            1.提高可伸缩性(Increase Scalability),
            2.可以利用JMS来缓解系统瓶颈(Reduce  Bottlenecks)
            3.提高系统对用户的响应能力
JMS公共API内部,和发送和接收JMS消息有关的JMS API接口常用的有7个:
        ·ConnectionFactory
        ·Destination
        ·Connection
        ·Session
        ·Message
        ·MessageProducer
        ·MessageConsumer
其中关键的5个部件,如图所示:
JMS API
查看大图请点击这里
Topic与Queue消息的区别
    Queue是一对一的消息传送,你可以看做是QQ应用程序中的一对一发送消息,一个人发出消息,只能有另外一个人阅读的到,中间需要通过一个队列支持Queue消息发送完毕后会保存在 JMS服务器的队列中,如果接收端接收以后将从队列中摘除。如图所示
http://l99raa.bay.livefilestore.com/y1pBMmD5G5XKCl5OM-cJV2PLvaN1GHnZ6mEwoC-DavOzzhHpFTnVeDERKWTIIW0eo3y9jxVmht3LGAYFS0hTvUIfEcIibsDOv82/p2p.png?psid=1
Queue消息发送到服务器,接收端会平均接收到发送端发送过来的消息,如图所示,一次发送150个消息,3个接收端每个收到50个
http://l99raa.bay.livefilestore.com/y1pDIUvLhMqs4VBwqyKXpE-4fImDxWL7hlUZUlqZYXRGbgLB81FDKJ_ejSl3eaFvP3M1FtV6y-p9KtxilPG7DO644DqjsqOzgrg/JMS-Dispatch.png

Topic  是一对多的消息传送,你可以看做是QQ应用程序 QQ群 聊天中的一对多发送消息,一个人发出消息,可以有多个订阅的人阅读的到,需要有一个消息主题作为支柱,Tocip消息发送完毕以后 无论有没有客户端接收,JMS服务器中的Topic消息都不会存在JMS服务器中。如图所示
http://l99raa.bay.livefilestore.com/y1p8wAJsZVtvdi48tGMpXfevLlrUzeeGXQwQJtjvsjJkcu4i4W8pzi9XuISR5pAMGZQU_k0xjSrLfFdJ0IBLP5ZjGVCIZJBv-ox/topic.png?psid=1
发送Topic消息无论连接在JMS服务器上有多少个接收端,将会收到同样的消息,而且Topic发送完毕以后不会保留在JMS服务器,否则将会和Topic消息的设计思想相互冲突。
http://l99raa.bay.livefilestore.com/y1plVf5LCPVGFnz7jN4CqQcianHZtJRkWOpPKaMLu12T56ZHANceuj3Ba5FcIj12H4Zo58X-o5ENe9ytURFpVvq5MJl27HzruBy/JMS-Dispatch-topic.png?psid=1

JMS消息的属性
JMS的每条消息分为 Header、Properties   、Body   3个属性 :
Header、Properties 中包含可设置的参数,分别是:
JMSDestination   消息发送的目的地
JMSDeliveryMode  传递模式, 有两种模式: PERSISTENT 和NON_PERSISTENT,PERSISTENT 表示该消息一定要被送到目的地,否则会导致应用错误。NON_PERSISTENT 表示偶然丢失该消息是被允许的,这两种模式使开发者可以在消息传递的可靠性和吞吐量之间找到平衡点。
JMSMessageID  唯一识别每个消息的标识,由JMS Provider 产生。
JMSTimestamp  一个消息被提交给JMS Provider 到消息被发出的时间。
JMSCorrelationID  用来连接到另外一个消息,典型的应用是在回复消息中连接到原消息。
JMSReplyTo        提供本消息回复消息的目的地址
JMSRedelivered  如果一个客户端收到一个设置了JMSRedelivered 属性的消息,则表示可能该客户端曾经在早些时候收到过该消息,但并没有签收(acknowledged)。
JMSType  消息类型的识别符。
JMSExpiration   消息过期时间,等于QueueSender 的send 方法中的timeToLive 值或TopicPublisher 的publish 方法中的timeToLive 值加上发送时刻的GMT 时间值。如果timeToLive值等于零,则JMSExpiration 被设为零,表示该消息永不过期。如果发送后,在消息过期时间之后消息还没有被发送到目的地,则该消息被清除。
JMSPriority  消息优先级,从0-9 十个级别,0-4 是普通消息,5-9 是加急消息。JMS 不要求JMS Provider 严格按照这十个优先级发送消息,但必须保证加急消息要先于普通消息到达。

Body 分为以下5种:
TextMessage
    这种类型携带了一个java.lang.String作为有效负载。它可以用于简单的文本消息交换,还可以用于更复杂的字符数据交换,比如XML文档等。
ObjectMessage
    这种类型携带了一个可序列化Java对象作为有效负载。它可以用于Java对象交换。
BytesMessage
    这种类型携带了一组原始类型字节流(primitive byte)作为有效负载。它可以使用应用程序的本机格式(native format)来交换数据,这种格式可能不兼容其他现有的Message类型。当JMS纯粹用于两个系统之间的消息传送时,也可以使用这种类型,而且该消 息的有效负载对JMS客户端来说是不透明的。
StreamMessage
    这种类型携带了一个Java原始数据类型流(int、double、char等)作为有效负载。它提供了一套将格式化字节流映射为Java原始数据类型的简便方法。在以固定顺序进行原始应用数据交换时,这种模型非常易于编程实现。
MapMessage
    这种类型携带了一组名/值对(name-value pair)作为有效负载。有效负载类似于一个java.util.Properties对象,除了有效负载值必须是Java原始数据类型或它们的包装器之外。MapMessage可以用于传送键入的数据。
如图所示:
jms的组成

JMS的消息存储
   
基于内存       最快
   基于文件       最慢
   基于数据库   比较慢
 
JMS的消息过滤功能
   这里selector是一个字符串,用来过滤消息。也就是说,这种方式可 以创建一个可以只接收特定消息的一个消费者。Selector的格式是类似于SQL-92的一种语法。可以用来比较消息头信息和属性。例如:发送消息属性中指定ip地址位数是奇数的接收,而ip地址位数是偶数的绝对不会接收到消息,具体实现方法请Google 关键字“jms selector
JMS消息的顺序
   乱序接收,顺序处理。也就是说,消息在发送,传输和接收过程中可能是乱序的,但消费者在接收到消息之后,并不立即处理,而是先将消息排序,然后在处理。 JMS消息头部的 JMSCorrelationID可以帮助我们完成这个工作。JMSCorrelationID存放了另一个消息的id。消息的发送者,如果要保证消息的顺序性,要将后发送的消息的JMSCorrelationID设置成前一个消息的id。消费者接收消息后,如果发现其头部有 JMSCorrelationID,则查看该消息是否已被处理过,如果没有,则等待该消息,至到该消息被处理后,才处理这个消息。这一工作需要发送者和接收者都记住已经发送和接收过的消息,以便于给后来的消息参考。由此可见,在JMS中设置JMS的head还是能起到不少作用的。

口水:
本文图绘采用国产软件 《亿图图示专家》所绘制,感谢《亿图图示专家》的所有开发人员制造出这么好的工具。

Java中间件:淘宝网系统高性能利器

【TechTarget中国原创】淘宝网是亚太最大的网络零售商圈,其知名度毋庸置疑,吸引着越来越多的消费者从街头移步这里,成为其忠实粉丝。如此多的用户和交易量,也意味着海量的信息处理,其背后的IT架构的稳定性、可靠性也显得尤为重要。那么,他们是怎么办到的呢?
  曾宪杰(花名花黎)是淘宝Java中间件团队成员,他认为大型网站就是要同时满足高访问量和高数据量的要求,核心是通过分布式系统解决数据的处理、存储及访问问题。
  消息中间件Notify
  早期,淘宝并没有Java中间件,其系统框架比较简单。下面我们就来看看Java中间件在淘宝的诞生和发展。首先要说的是实现系统松耦合和异步处理的消息中间件Notify,这是一个高性能、高可靠、可扩展组件,轻量级支持最终一致性和订阅者集群。所谓订阅者集群,即将订阅消息的客户端分为多个集群,集群之间采用Topic方式,让每个集群都能收到消息,集群之中再按照Queue的方式,仅由一个客户端来处理消息。
  对于淘宝来说,最终一致性至关重要。有过淘宝经验的人都知道,在我们完成付款之后,订单状态会立刻更改为已付款。如果用户付款之后,淘宝不能通知支付宝为该用户的账号充值,商家也不能知道用户已经付款,也就是整个交易的所有参与方不能实现最终状态一致性的话,整个交易也就无法继续下去。曾宪杰笑言:“如果真的发生这样的情况,那么淘宝就不用做了。”
  在实现消息的可靠性上,淘宝采用Oracle+小型机(IBM)+高端存储(EMC)的形式,写双份Mysql,同时基于文件和内存。Notify目前每天消息总量达到4.4亿,每天消息投递条次约15亿次,总共有78个消息主题,786种消息类型,部分消息订阅者超过30个集群。下图是淘宝在应用了Notify之后的系统架构图:
Notify
  淘宝服务框架——HSF
  应用了消息中间件之后,淘宝仍旧面临着一系列问题,比如上百人维护一个代码百万行的前台核心应用;多个业务系统中的代码重复编写以及数据库连接数接近瓶颈。那怎么解决呢?服务化成为淘宝的上选。应用服务化解决了业务核心的稳定和一致的问题,重要数据库的连接数也得到了缓解;此外,系统分解后,效率和稳定性也得到了显著提升。
  曾宪杰介绍他们的这个服务框架称之为HSF。目前HSF线上提供的服务数量超过六百个,每日的调用总量达到150亿以上,现在甚至更高。下图是应用了HSF之后的系统架构图:
服务框架HSF
  淘宝分布式数据层TDDL
  在淘宝的数据库架构演进过程中,为了更好地支持分库分表以及读写分离,进行了一定的封装。对上层应用而言仍旧操作JDBC,实则是在使用淘宝分布式数据层(TDDL),它能实现SQL解析、规则路由、数据合并;既可以用jar的方式在客户端直接连接数据库,也可以让客户端通过DBProxy服务器访问数据库;具备三层的数据源结构,还支持非对称数据复制。
  目前TDDL每日SQL执行量达到30亿以上,每日的数据复制量约为2.8亿多。下图是淘宝加上TDDL之后的系统架构:
TDDL
  尽管目前淘宝的Java中间件发展顺利,但也并不意味已经解决了一切问题。目前仍旧存在一些问题需要解决,曾宪杰表示在这些问题上,他们希望通过开源的途径得到解决,同时能够根据业务需求提供相应的新功能,另外系统的稳定性仍旧是他们要关注的内容。

2013年6月14日星期五

[翻译] [RabbitMQ+Python入门经典] 兔子和兔子窝

 RabbitMQ作为一个工业级的消息队列服务器,在其客户端手册列表的Python段当中推荐了一篇blog ,作为RabbitMQ+Python的入门手册再合适不过了。不过,正如其标题Rabbit and Warrens (兔子和养兔场)一样,,这篇英文写的相当俏皮,以至于对于我等非英文读者来说不像一般的技术文档那么好懂,所以,翻译一下吧。翻译过了,希望其他人可以少用一些时间。翻译水平有限,不可能像原文一样俏皮,部分地方可能就意译了,希望以容易懂为准。想看看老外的幽默的,推荐去看原文,其实,也不是那么难理解……
  兔子和兔子窝
  当时我们的动机很简单:从生产环境的电子邮件处理流程当中分支出一个特定的离线分析流程。我们开始用的MySQL,将要处理的东西放在表里面,另一个程序从中取。不过很快,这种设计的丑陋之处就显现出来了…… 你想要多个程序从一个队列当中取数据来处理?没问题,我们硬编码程序的个数好了……什么?还要能够允许程序动态地增加和减少的时候动态进行压力分配?
  是的,当年我们想的简单的东西(做一个分支处理)逐渐变成了一个棘手的问题。以前拿着锤子(MySQL)看所有东西都是钉子(表)的年代是多么美好……
  在搜索了一下之后,我们走进了消息队列(message queue)的大门。不不,我们当然知道消息队列是什么,我们可是以做电子邮件程序谋生的。我们实现过各种各样的专业的,高速的内存队列用来做电子邮件处理。我们不知道的是那一大类现成的、通用的消息队列(MQ)服务器——无论是用什么语言写出的,不需要复杂的装配的,可以自然的在网络上的应用程序之间传送数据的一类程序。不用我们自己写?看看再说。
  让大家看看你们的Queue吧……
  过去的4年里,人们写了有好多好多的开源的MQ服务器啊。其中大多数都是某公司例如LiveJournal写出来用来解决特定问题的。它们的确不关心上面跑的是什么类型的消息,不过他们的设计思想通常是和创建者息息相关的(消息的持久化,崩溃恢复等通常不在他们考虑范围内)。不过,有三个专门设计用来做及其灵活的消息队列的程序值得关注:
  Apache ActiveMQ
  ZeroMQ
  RabbitMQ
  Apache ActiveMQ曝光率最高,不过看起来它有些问题,可能会造成丢消息。不可接受,下一个。
  ZeroMQ和RabbitMQ都支持一个开源的消息协议,成为AMQP。AMQP的一个优点是它是一个灵活和开放的协议,以便和另外两个商业化的Message Queue (IBM和Tibco)竞争,很好。不过ZeroMQ不支持消息持久化和崩溃恢复,不太好。剩下的只有RabbitMQ了。如果你不在意消息持久化和崩溃恢复,试试ZeroMQ吧,延迟很低,而且支持灵活的拓扑。
  剩下的只有这个吃胡萝卜的家伙了……

 
  当我读到它是用Erlang写的时候,RabbitMQ震了我一下。Erlang 是爱立信开发的高度并行的语言,用来跑在电话交换机上。是的,那些要求6个9的在线时间的东西。在Erlang当中,充斥着大量轻量进程,它们之间用消息传递来通信。听起来思路和我们用消息队列的思路是一样的,不是么?
  而且,RabbitMQ支持持久化。是的,如果RabbitMQ死掉了,消息并不会丢失,当队列重启,一切都会回来。而且,正如在 DigiTar(注:原文作者的公司)做事情期望的那样,它可以和Python无缝结合 。除此之外,RabbitMQ的文档相当的……恐怖。如果你懂AMQP,这些文档还好,但是有多少人懂AMQP?这些文档就像MySQL的文档假设你已经懂了SQL一样……不过没关系啦。
  好了,废话少说。这里是花了一周时间阅读关于AMQP和关于它如何在RabbitMQ上工作的文档之后的一个总结,还有,怎么在Python当中使用。
  开始吧
  AMQP当中有四个概念非常重要:虚拟主机(virtual host),交换机(exchange),队列(queue)和绑定(binding)。一个虚拟主机持有一组交换机、队列和绑定。为什么需要多个虚拟主机呢?很简单,RabbitMQ当中,用户只能在虚拟主机的粒度进行权限控制。因此,如果需要禁止A组访问B组的交换机/队列/绑定,必须为A和B分别创建一个虚拟主机。每一个RabbitMQ服务器都有一个默认的虚拟主机“/”。如果这就够了,那现在就可以开始了。
  交换机,队列,还有绑定……天哪!
  刚开始我思维的列车就是在这里脱轨的…… 这些鬼东西怎么结合起来的?
  队列(Queues)是你的消息(messages)的终点,可以理解成装消息的容器。消息就一直在里面,直到有客户端(也就是消费者,Consumer)连接到这个队列并且将其取走为止。不过。你可以将一个队列配置成这样的:一旦消息进入这个队列,biu~,它就烟消云散了。这个有点跑题了……
  需要记住的是,队列是由消费者(Consumer)通过程序建立的,不是通过配置文件或者命令行工具。这没什么问题,如果一个消费者试图创建一个已经存在的队列,RabbitMQ就会起来拍拍他的脑袋,笑一笑,然后忽略这个请求。因此你可以将消息队列的配置写在应用程序的代码里面。这个概念不错。
  OK,你已经创建并且连接到了你的队列,你的消费者程序正在百无聊赖的敲着手指等待消息的到来,敲啊,敲啊…… 没有消息。发生了什么?你当然需要先把一个消息放进队列才行。不过要做这个,你需要一个交换机(Exchange)……
  交换机可以理解成具有路由表的路由程序,仅此而已。每个消息都有一个称为路由键(routing key)的属性,就是一个简单的字符串。交换机当中有一系列的绑定(binding),即路由规则(routes),例如,指明具有路由键 “X” 的消息要到名为timbuku的队列当中去。先不讨论这个,我们有点超前了。
  你的消费者程序要负责创建你的交换机们 (复数)。啥?你是说你可以有多个交换机?是的,这个可以有,不过为啥?很简单,每个交换机在自己独立的进程当中执行,因此增加多个交换机就是增加多个进程,可以充分利用服务器上的CPU核以便达到更高的效率。例如,在一个8 核的服务器上,可以创建5个交换机来用5个核,另外3个核留下来做消息处理。类似的,在RabbitMQ的集群当中,你可以用类似的思路来扩展交换机一边获取更高的吞吐量。
  OK,你已经创建了一个交换机。但是他并不知道要把消息送到哪个队列。你需要路由规则,即绑定(binding)。一个绑定就是一个类似这样的规则:将交换机“desert(沙漠)”当中具有路由键“阿里巴巴”的消息送到队列“hideout(山洞)”里面去。换句话说,一个绑定就是一个基于路由键将交换机和队列连接起来的路由规则。例如,具有路由键“audit”的消息需要被送到两个队列,“log-forever”和“alert-the- big-dude”。要做到这个,就需要创建两个绑定,每个都连接一个交换机和一个队列,两者都是由“audit”路由键触发。在这种情况下,交换机会复制一份消息并且把它们分别发送到两个队列当中。交换机不过就是一个由绑定构成的路由表。
  现在复杂的东西来了:交换机有多种类型。他们都是做路由的,不过接受不同类型的绑定。为什么不创建一种交换机来处理所有类型的路由规则呢?因为每种规则用来做匹配分子的CPU开销是不同的。例如,一个“topic”类型的交换机试图将消息的路由键与类似“dogs.* ” 的模式进行匹配。匹配这种末端的通配符比直接将路由键与“dogs ”比较(“direct”类型的交换机)要消耗更多的CPU。如果你不需要“topic”类型的交换机带来的灵活性,你可以通过使用“direct”类型的交换机获取更高的处理效率。那么有哪些类型,他们又是怎么处理的呢?
  Fanout Exchange——不处理路由键。你只需要简单的将队列绑定到交换机上。一个发送到交换机的消息都会被转发到与该交换机绑定的所有队列上。很像子网广播,每台子网内的主机都获得了一份复制的消息。Fanout交换机转发消息是最快的。
  Direct Exchange——处理路由键。需要将一个队列绑定到交换机上,要求该消息与一个特定的路由键完全匹配。这是一个完整的匹配。如果一个队列绑定到该交换机上要求路由键 “dog”,则只有被标记为“dog ”的消息才被转发,不会转发dog.puppy ,也不会转发dog.guard ,只会转发dog 。
  Topic Exchange——将路由键和某模式进行匹配。此时队列需要绑定要一个模式上。符号“#”匹配一个或多个词,符号“*”匹配不多不少一个词。因此“audit.#”能够匹配到“audit.irs.corporate ”,但是“audit.* ” 只会匹配到“audit.irs ”。
 持久化这些小东西们
  你花了大量的时间来创建队列、交换机和绑定,然后,砰~服务器程序挂了。你的队列、交换机和绑定怎么样了?还有,放在队列里面但是尚未处理的消息们呢?
  放松~如果你是用默认参数构造的这一切的话,那么,他们,都,biu~,灰飞烟灭了。是的,RabbitMQ重启之后会干净的像个新生儿。你必须重做所有的一切,亡羊补牢,如何避免将来再度发生此类杯具?
  队列和交换机有一个创建时候指定的标志durable,直译叫做坚固的。durable的唯一含义就是具有这个标志的队列和交换机会在重启之后重新建立,它不表示说在队列当中的消息会在重启后恢复。那么如何才能做到不只是队列和交换机,还有消息都是持久的呢?
  但是首先一个问题是,你真的需要消息是持久的吗?对于一个需要在重启之后回复的消息来说,它需要被写入到磁盘上,而即使是最简单的磁盘操作也是要消耗时间的。如果和消息的内容相比,你更看重的是消息处理的速度,那么不要使用持久化的消息。不过对于我们@DigiTar来说,持久化很重要。
  当你将消息发布到交换机的时候,可以指定一个标志“Delivery Mode”(投递模式)。根据你使用的AMQP的库不同,指定这个标志的方法可能不太一样(我们后面会讨论如何用Python搞定)。简单的说,就是将 Delivery Mode设置成2,也就是持久的(persistent)即可。一般的AMQP库都是将Delivery Mode设置成1,也就是非持久的。所以要持久化消息的步骤如下:
  •   将交换机设成 durable。
  •   将队列设成 durable。
  •   将消息的 Delivery Mode 设置成2 。
  就这样,不是很复杂,起码没有造火箭复杂,不过也有可能犯点小错误。
  下面还要罗嗦一个东西……绑定(Bindings)怎么办?我们无法在创建绑定的时候设置成durable。没问题,如果你绑定了一个 durable的队列和一个durable的交换机,RabbitMQ会自动保留这个绑定。类似的,如果删除了某个队列或交换机(无论是不是 durable),依赖它的绑定都会自动删除。
  注意两点:
  RabbitMQ不允许你绑定一个非坚固(non-durable)的交换机和一个durable的队列。反之亦然。要想成功必须队列和交换机都是durable的。 一旦创建了队列和交换机,就不能修改其标志了。例如,如果创建了一个non-durable的队列,然后想把它改变成durable的,唯一的办法就是删除这个队列然后重现创建。因此,最好仔细检查创建的标志。
  开始喂蛇了~【译注】说喂蛇是因为Python的图标是条蛇。
  AMQP的一个空白地带是如何在Python当中使用。对于其他语言有一大坨材料。
以下是引用片段:
Java – http://www.rabbitmq.com/java-client.html  
.NET – http://www.rabbitmq.com/releases/rabbitmq-dotnet-client/v1.5.0/rabbitmq-dotnet-client-1.5.0-user-guide.pdf  
Ruby – http://somic.org/2008/06/24/ruby-amqp-rabbitmq-example/ 
  但是对Python老兄来说,你需要花点时间来挖掘一下。所以我写了这个,这样别的家伙们就不需要经历我这种抓狂的过程了。
  首先,我们需要一个Python的AMQP库。有两个可选:
  •   py-amqplib——通用的AMQP
  •   txAMQP——使用 Twisted 框架的AMQP库,因此允许异步I/O。
  根据你的需求,py-amqplib或者txAMQP都是可以的。因为是基于Twisted的,txAMQP可以保证用异步IO构建超高性能的 AMQP程序。但是Twisted编程本身就是一个很大的主题……因此清晰起见,我们打算用 py-amqplib。更新:请参见 Esteve Fernandez关于txAMQP的使用和代码样例的回复 。
  AMQP支持在一个TCP连接上启用多个MQ通信channel,每个channel都可以被应用作为通信流。每个AMQP程序至少要有一个连接和一个channel。
以下是引用片段:
view plaincopy to clipboardprint? 
<SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">from</SPAN>    
 amqplib <SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">import</SPAN>   
 client_0_8 <SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">as</SPAN>    
 amqp<BR>    
conn = amqp.<SPAN style="COLOR: black">Connection</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
host=<SPAN style="COLOR: #483d8b">"localhost:5672 "</SPAN>    
, userid=<SPAN style="COLOR: #483d8b">"guest"</SPAN>    
,<BR>    
password=<SPAN style="COLOR: #483d8b">"guest"</SPAN>    
, virtual_host=<SPAN style="COLOR: #483d8b">"/"</SPAN>    
, insist=<SPAN style="COLOR: #008000">False</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
chan = conn.<SPAN style="COLOR: black">channel</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: black">)</SPAN>   
from 
 amqplib import 
 client_0_8 as 
 amqp 
conn = amqp.Connection 

host="localhost:5672 " 
, userid="guest" 

password="guest" 
, virtual_host="/" 
, insist=False 

chan = conn.channel 

)
  每个channel都被分配了一个整数标识,自动由Connection()类的.channel()方法维护。或者,你可以使用.channel(x)来指定channel标识,其中x是你想要使用的channel标识。通常情况下,推荐使用.channel()方法来自动分配 channel标识,以便防止冲突。
  现在我们已经有了一个可以用的连接和channel。现在,我们的代码将分成两个应用,生产者(producer)和消费者(consumer)。我们先创建一个消费者程序,他会创建一个叫做“po_box”的队列和一个叫“sorting_room”的交换机:
以下是引用片段:
view plaincopy to clipboardprint? 
chan.<SPAN style="COLOR: black">queue_declare</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
queue=<SPAN style="COLOR: #483d8b">"po_box"</SPAN>    
, durable=<SPAN style="COLOR: #008000">True</SPAN>    
, exclusive=<SPAN style="COLOR: #008000">False</SPAN>    
, auto_delete=<SPAN style="COLOR: #008000">False</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
chan.<SPAN style="COLOR: black">exchange_declare</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
exchange=<SPAN style="COLOR: #483d8b">"sorting_room"</SPAN>    
, <SPAN style="COLOR: #008000">type</SPAN>    
=<SPAN style="COLOR: #483d8b">"direct"</SPAN>    
, durable=<SPAN style="COLOR: #008000">True</SPAN>    
, auto_delete=<SPAN style="COLOR: #008000">False</SPAN>    
,<SPAN style="COLOR: black">)</SPAN>   
chan.queue_declare 

queue="po_box" 
, durable=True 
, exclusive=False 
, auto_delete=False 

chan.exchange_declare 

exchange="sorting_room" 
, type 
="direct" 
, durable=True 
, auto_delete=False 
,)
  这段代码干了啥?首先,它创建了一个名叫“po_box ”的队列,它是durable的(重启之后会重新建立),并且最后一个消费者断开的时候不会自动删除(auto_delete=False )。在创建durable的队列(或者交换机)的时候,将auto_delete设置成false是很重要的,否则队列将会在最后一个消费者断开的时候消失,与durable与否无关。如果将durable和auto_delete都设置成True,只有尚有消费者活动的队列可以在RabbitMQ意外崩溃的时候自动恢复。
  (你可以注意到了另一个标志,称为“exclusive”。如果设置成True,只有创建这个队列的消费者程序才允许连接到该队列。这种队列对于这个消费者程序是私有的)。
  还有另一个交换机声明,创建了一个名字叫“sorting_room”的交换机。auto_delete和durable的含义和队列是一样的。但是,.excange_declare() 还有另外一个参数叫做type,用来指定要创建的交换机的类型(如前面列出的): fanout , direct 和 topic .
  到此为止,你已经有了一个可以接收消息的队列和一个可以发送消息的交换机。不过我们需要创建一个绑定,把它们连接起来。
  chan.queue_bind(queue=”po_box”, exchange=”sorting_room”, routing_key=”jason”)这个绑定的过程非常直接。任何送到交换机“sorting_room ”的具有路由键“jason ” 的消息都被路由到名为“po_box ” 的队列。
  现在,你有两种方法从队列当中取出消息。第一个是调用chan.basic_get() ,主动从队列当中拉出下一个消息(如果队列当中没有消息,chan.basic_get()会返回None, 因此下面代码当中print msg.body 会在没有消息的时候崩掉):
以下是引用片段:
view plaincopy to clipboardprint? 
msg = chan.<SPAN style="COLOR: black">basic_get</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: #483d8b">"po_box"</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
<SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">print</SPAN>    
 msg.<SPAN style="COLOR: black">body</SPAN>    
<BR>    
chan.<SPAN style="COLOR: black">basic_ack</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
msg.<SPAN style="COLOR: black">delivery_tag</SPAN>    
<SPAN style="COLOR: black">)</SPAN>   
msg = chan.basic_get 

"po_box" 

print 
 msg.body 
chan.basic_ack 

msg.delivery_tag 
)
  但是如果你想要应用程序在消息到达的时候立即得到通知怎么办?这种情况下不能使用chan.basic_get() ,你需要用chan.basic_consume() 注册一个新消息到达的回调。
以下是引用片段:
view plaincopy to clipboardprint? 
<SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">def</SPAN>    
 recv_callback<SPAN style="COLOR: black">(</SPAN>    
msg<SPAN style="COLOR: black">)</SPAN>    
:<BR>    
    <SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">print</SPAN>    
 <SPAN style="COLOR: #483d8b">'Received: '</SPAN>    
 + msg.<SPAN style="COLOR: black">body</SPAN>    
<BR>    
chan.<SPAN style="COLOR: black">basic_consume</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
queue=<SPAN style="COLOR: #483d8b">'po_box'</SPAN>    
, no_ack=<SPAN style="COLOR: #008000">True</SPAN>    
,<BR>    
callback=recv_callback, consumer_tag=<SPAN style="COLOR: #483d8b">"testtag"</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
<SPAN style="FONT-WEIGHT: bold; COLOR: #ff7700">while</SPAN>    
 <SPAN style="COLOR: #008000">True</SPAN>    
:<BR>    
    chan.<SPAN style="COLOR: black">wait</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
chan.<SPAN style="COLOR: black">basic_cancel</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: #483d8b">"testtag"</SPAN>    
<SPAN style="COLOR: black">)</SPAN>   
def 
 recv_callback( 
msg) 

    print 
 'Received: ' 
 + msg.body 
chan.basic_consume 

queue='po_box' 
, no_ack=True 

callback=recv_callback, consumer_tag="testtag" 

while 
 True 

    chan.wait 


chan.basic_cancel 

"testtag" 
)
  chan.wait()放在一个无限循环里面,这个函数会等待在队列上,直到下一个消息到达队列。chan.basic_cancel()用来注销该回调函数。参数consumer_tag 当中指定的字符串和chan.basic_consume()注册的一直。在这个例子当中chan.basic_cancel()不会被调用到,因为上面是个无限循环…… 不过你需要知道这个调用,所以我把它放在了代码里。
  需要注意的另一个东西是no_ack参数。这个参数可以传给chan.basic_get() 和chan.basic_consume(),默认是false。当从队列当中取出一个消息的时候,RabbitMQ需要应用显式地回馈说已经获取到了该消息。如果一段时间内不回馈,RabbitMQ 会将该消息重新分配给另外一个绑定在该队列上的消费者。另一种情况是消费者断开连接,但是获取到的消息没有回馈,则RabbitMQ同样重新分配。如果将no_ack 参数设置为true,则py-amqplib会为下一个AMQP请求添加一个no_ack属性,告诉AMQP服务器不需要等待回馈。但是,大多数时候,你也许想要自己手工发送回馈,例如,需要在回馈之前将消息存入数据库。回馈通常是通过调用chan.basic_ack() 方法,使用消息的delivery_tag 属性作为参数。参见chan.basic_get() 的实例代码。
  好了,这就是消费者的全部代码。(下载:amqp_consumer.py )
  不过没有人发送消息的话,要消费者何用?所以需要一个生产者。下面的代码示例表明如何将一个简单消息发送到交换区“sorting_room ”,并且标记为路由键“jason ” :
以下是引用片段:
view plaincopy to clipboardprint? 
msg = amqp.<SPAN style="COLOR: black">Message</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: #483d8b">"Test message!"</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
msg.<SPAN style="COLOR: black">properties</SPAN>    
<SPAN style="COLOR: black">[</SPAN>    
<SPAN style="COLOR: #483d8b">"delivery_mode"</SPAN>    
<SPAN style="COLOR: black">]</SPAN>    
 = <SPAN style="COLOR: #ff4500">2</SPAN>    
<BR>    
chan.<SPAN style="COLOR: black">basic_publish</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
msg,exchange=<SPAN style="COLOR: #483d8b">"sorting_room"</SPAN>    
,routing_key=<SPAN style="COLOR: #483d8b">"jason"</SPAN>    
<SPAN style="COLOR: black">)</SPAN>   
msg = amqp.Message 

"Test message!" 

msg.properties 

"delivery_mode" 

 = 2 
chan.basic_publish 

msg,exchange="sorting_room" 
,routing_key="jason" 
)
  你也许注意到我们设置消息的delivery_mode 属性为2,因为队列和交换机都设置为durable的,这个设置将保证消息能够持久化,也就是说,当它还没有送达消费者之前如果RabbitMQ重启则它能够被恢复。
  剩下的最后一件事情(生产者和消费者都需要调用的)是关闭channel和连接:
以下是引用片段:
view plaincopy to clipboardprint? 
chan.<SPAN style="COLOR: black">close</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: black">)</SPAN>    
<BR>    
conn.<SPAN style="COLOR: black">close</SPAN>    
<SPAN style="COLOR: black">(</SPAN>    
<SPAN style="COLOR: black">)</SPAN>   
chan.close 


conn.close 

)
  很简单吧。(下载:amqp_publisher.py )
  来真实地跑一下吧……
  现在我们已经写好了生产者和消费者,让他们跑起来吧。假设你的RabbitMQ在localhost上安装并且运行。
  打开一个终端,执行python ./amqp_consumer.py 让消费者运行,并且创建队列、交换机和绑定。
  然后在另一个终端运行python ./amqp_publisher.py “AMQP rocks.” 。如果一切良好,你应该能够在第一个终端看到输出的消息。
  付诸使用吧
  我知道这个教程是非常粗浅的关于AMQP/RabbitMQ和如何使用Python访问的教程。希望这个可以说明所有的概念如何在Python当中被组合起来。如果你发现任何错误,请联系原作者(williamsjj@digitar.com ) 【译注:如果是翻译问题请联系译者】。同时,我很高兴回答我知道的问题。【译注:译者也是一样的】。接下来是,集群化(clustering)!不过我需要先把它弄懂再说。

RabbitMQ持久化的PHP使用示例

      在使用RabbitMQ时候想对消息进行持久化,一开始exchange和queue做了持久化设置setFlags(AMQP_DURABLE),重启rabbitmq-server后发现数据丢了。后来看官方手册里的publish方法里有delivery_mode、priority两个属性,delivery_mode:1非持久化,delivery_mode:2持久化,其中priority是设置消息优先级的可设置0-9级别,但是官方解释好像并未实现,测试中发现并未起作用。即使设置了持久化,也不能完全保证消息不丢失,比如RabbitMQ接受到消息后还没来得及写到磁盘就宕机了,那么就会有小概率的信息丢失。

具体使用示例如下:
<?php
//连接RabbitMQ
$conn_args = array('host' =>'localhost', 'port'=> '5672', 'login' =>'guest', 'password'=> 'guest', 'vhost' => '/');
$conn = new AMQPConnection($conn_args);
$conn->connect();

$channel = new AMQPChannel($conn);
$exchangeName = 'reported';
$queueName = 'request';
$queueKey = 'request_item';

if( isset($_GET['read']) ){ //读队列

$q = new AMQPQueue($channel);
$q->setName($queueName);
$q->setFlags(AMQP_DURABLE);
$messages = '';
try {
$q->bind($exchangeName, $queueKey);
$messages = $q->get(AMQP_AUTOACK); //消息获取
!empty($messages) && $messages = json_decode($messages->getBody(), true);
} catch (Exception $e) {
var_dump($e);
}
var_dump($messages);
$conn->disconnect();

}else{ //写队列

$ex = new AMQPExchange($channel);
$ex->setName($exchangeName);
$ex->setType(AMQP_EX_TYPE_TOPIC);
$ex->setFlags(AMQP_DURABLE); //exchange持久化
$ex->declare();

$q = new AMQPQueue($channel);
$q->setName($queueName);
$q->setFlags(AMQP_DURABLE); //queue持久化
$q->declare();
$q->bind($exchangeName, $queueKey);
$channel->startTransaction();
$message = array('content' => 'test','time' => time());

/**
* 消息持久化,delivery_mode:2持久化、delivery_mode:1非持久化,其中priority是设置消息的优先级,测试中发现并未起作用。
* 消息还有其他属性,请参考http://www.php.net/manual/zh/amqpexchange.publish.php
*/
$result = $ex->publish(json_encode($message), $queueKey, AMQP_NOPARAM, array('delivery_mode'=>2, 'priority'=> 9));
var_dump($result);
$channel->commitTransaction();
$conn->disconnect();
}

?>

【架构】关于RabbitMQ

1      什么是RabbitMQ

RabbitMQ是实现AMQP(高级消息队列协议)的消息中间件的一种,最初起源于金融系统,用于在分布式系统中存储转发消息,在易用性、扩展性、高可用性等方面表现不俗。消息中间件主要用于组件之间的解耦,消息的发送者无需知道消息使用者的存在,反之亦然:

 
单向解耦

 
双向解耦(如:RPC
    例如一个日志系统,很容易使用RabbitMQ简化工作量,一个Consumer可以进行消息的正常处理,另一个Consumer负责对消息进行日志记录,只要在程序中指定两个Consumer所监听的queue以相同的方式绑定到同一个exchange即可,剩下的消息分发工作由RabbitMQ完成。
 

使用RabbitMQ server需要:
1. ErLang语言包;
2. RabbitMQ安装包;
RabbitMQ同时提供了java的客户端(一个jar包)。

2      概念和特性

2.1      交换机(exchange):

1. 接收消息,转发消息到绑定的队列。四种类型:direct, topic, headers and fanout
direct:转发消息到routigKey指定的队列
topic:按规则转发消息(最灵活)
headers:(这个还没有接触到)
fanout:转发消息到所有绑定队列
2. 如果没有队列绑定在交换机上,则发送到该交换机上的消息会丢失。
3. 一个交换机可以绑定多个队列,一个队列可以被多个交换机绑定。
4. topic类型交换器通过模式匹配分析消息的routing-key属性。它将routing-keybinding-key的字符串切分成单词。这些单词之间用点隔开。它同样也会识别两个通配符:#匹配0个或者多个单词,*匹配一个单词。例如,binding key*.stock.#匹配routing keyusd.stcokeur.stock.db,但是不匹配stock.nana
还有一些其他的交换器类型,如headerfailoversystem等,现在在当前的RabbitMQ版本中均未实现。
5. 因为交换器是命名实体,声明一个已经存在的交换器,但是试图赋予不同类型是会导致错误。客户端需要删除这个已经存在的交换器,然后重新声明并且赋予新的类型。
6. 交换器的属性:
持久性:如果启用,交换器将会在server重启前都有效。
自动删除:如果启用,那么交换器将会在其绑定的队列都被删除掉之后自动删除掉自身。
惰性:如果没有声明交换器,那么在执行到使用的时候会导致异常,并不会主动声明。

2.2      队列(queue):

1. 队列是RabbitMQ内部对象,存储消息。相同属性的queue可以重复定义。
2. 临时队列。channel.queueDeclare(),有时不需要指定队列的名字,并希望断开连接时删除队列。
3. 队列的属性:
持久性:如果启用,队列将会在server重启前都有效。
自动删除:如果启用,那么队列将会在所有的消费者停止使用之后自动删除掉自身。
惰性:如果没有声明队列,那么在执行到使用的时候会导致异常,并不会主动声明。
排他性:如果启用,队列只能被声明它的消费者使用。
这些性质可以用来创建例如排他和自删除的transient或者私有队列。这种队列将会在所有链接到它的客户端断开连接之后被自动删除掉。它们只是短暂地连接到server,但是可以用于实现例如RPC或者在AMQ上的对等通信。4. RPC的使用是这样的:RPC客户端声明一个回复队列,唯一命名(例如用UUID),并且是自删除和排他的。然后它发送请求给一些交换器,在消息的reply-to字段中包含了之前声明的回复队列的名字。RPC服务器将会回答这些请求,使用消息的reply-to作为routing key(默认绑定器会绑定所有的队列到默认交换器,名称为“amp.交换器类型名”)发送到默认交换器。注意这仅仅是惯例而已,可以根据和RPC服务器的约定,它可以解释消息的任何属性(甚至数据体)来决定回复给谁。

2.3      消息传递:

1. 消息在队列中保存,以轮询的方式将消息发送给监听消息队列的消费者,可以动态的增加消费者以提高消息的处理能力。
2. 为了实现负载均衡,可以在消费者端通知RabbitMQ,一个消息处理完之后才会接受下一个消息。
channel.basic_qos(prefetch_count=1)
注意:要防止如果所有的消费者都在处理中,则队列中的消息会累积的情况。
3. 消息有14个属性,最常用的几种:
deliveryMode:持久化属性
contentType:编码
replyTo:指定一个回调队列
correlationId:消息id
实例代码:
4. 消息生产者可以选择是否在消息被发送到交换器并且还未投递到队列(没有绑定器存在)和/或没有消费者能够立即处理的时候得到通知。通过设置消息的mandatory/immediate属性为真,这些投递保障机制的能力得到了强化。
5. 此外,一个生产者可以设置消息的persistent属性为真。这样一来,server将会尝试将这些消息存储在一个稳定的位置,直到server崩溃。当然,这些消息肯定不会被投递到非持久的队列中。

2.4      高可用性(HA):

1. 消息ACK,通知RabbitMQ消息已被处理,可以从内存删除。如果消费者因宕机或链接失败等原因没有发送ACK(不同于ActiveMQ,在RabbitMQ里,消息没有过期的概念),则RabbitMQ会将消息重新发送给其他监听在队列的下一个消费者。
channel.basicConsume(queuename, noAck=false, consumer);
2. 消息和队列的持久化。定义队列时可以指定队列的持久化属性(问:持久化队列如何删除?)
channel.queueDeclare(queuename, durable=true, false, false, null);
发送消息时可以指定消息持久化属性:
channel.basicPublish(exchangeName, routingKey,
            MessageProperties.PERSISTENT_TEXT_PLAIN,
            message.getBytes());
这样,即使RabbitMQ服务器重启,也不会丢失队列和消息。
3. publisher confirms
4. master/slave机制,配合Mirrored Queue,这种情况下,publisher会正常发送消息和接收消息的confirm,但对于subscriber来说,需要接收Consumer Cancellation Notifications来得到主节点失败的通知,然后re-consume from the queue,此时要求client有处理重复消息的能力。注意:如果queue在一个新加入的节点上增加了一个slave,此时slave上没有此前queue的信息(目前还没有同步机制)。
(通过命令行或管理插件可以查看哪个slave是同步的:
rabbitmqctl list_queues name slave_pids synchronised_slave_pids
    当一个slave重新加入mirrored-queue时,如果queuedurable的,则会被清空。

2.5      集群(cluster):

1. 不支持跨网段(如需支持,需要shovelfederation插件)
2. 可以随意的动态增加或减少、启动或停止节点,允许节点故障
3. 集群分为RAM节点和DISK节点,一个集群最好至少有一个DISK节点保存集群的状态。
4. 集群的配置可以通过命令行,也可以通过配置文件,命令行优先。

3      使用

3.1      简易使用流程
 

3.2      RabbitMQOpenStack中的使用

 


    Openstack中,组件之间对RabbitMQ使用基本都是“Remote Procedure Calls”的方式。每一个Nova服务(比如计算服务、存储服务等)初始化时会创建两个队列,一个名为“NODE-TYPE.NODE-ID”,另一个名为“NODE-TYPE”,NODE-TYPE是指服务的类型,NODE-ID指节点名称。
    从抽象层面上讲,RabbitMQ的组件的使用类似于下图所示:
每个服务会绑定两个队列到同一个topic类型的exchange,从不同的队列中接收不同类型的消息。消息的发送者如果关心消息的返回值,则会监听另一个队列,该队列绑定在一个direct类型的exchange。接受者收到消息并处理后,会将消息的返回发送到此exchange
Openstack中,如果不关心消息返回,消息的流程图如下:

 
    如果关心消息返回值,流程图如下:


 

3.3      为什么要使用RabbitMQ

曾经有过一个人做过一个测试(http://www.cnblogs.com/amityat/archive/2011/08/31/2160293.html),发送1百万个并发消息,对性能有很高的需求,于是作者对比了RabbitMQMSMQActiveMQZeroMQueue,整个过程共产生1百万条1K的消息。测试的执行是在一个Windows Vista上进行的,测试结果如下:

    虽然ZeroMQ性能较高,但这个产品不提供消息持久化,需要自己实现审计和数据恢复,因此在易用性和HA上不是令人满意,通过测试结果可以看到,RabbitMQ的性能确实不错。
    我在本机也做了一些测试,但我的测试是基于组件的原生配置,没有做任何的配置优化,因此总觉的不靠谱。我只测试了RabbitMQActiveMQ两款产品,虽然网上都说ActiveMQ性能不如前者,但平心而论,ActiveMQ提供了很多配置,存在很大的调优空间,也许修改一个配置参数就会使组件的性能有一个质的飞跃。