Python通过amqp消息队列协议中的Qpid实现数据通信

简介:

简介:

    这两天看了消息队列通信,打算在配置平台上应用起来。以前用过zeromq但是这东西太快了,还有就是rabbitmq有点大,新浪的朋友推荐了qpid,简单轻便。自己总结了下文档,大家可以瞅瞅。


AMQP(消息队列协议Advanced Message Queuing Protocol)是一种消息协议 ,等同于JMS,但是JMS只是java平台的方案,AMQP是一个跨语言的协议。


AMQP 不分语言平台,主流的语言都支持,运维这边的perl,python,ruby更是支持,所以大家就放心用吧。


主流的消息队列通信类型:

1
2
3
4
5
6
7
点对点:A 发消息给 B。
广播:A 发给所有其他人的消息
组播:A 发给多个但不是所有其他人的消息。
Requester/response:类似访问网页的通信方式,客户端发请求并等待,服务端回复该请求
Pub-sub:类似杂志发行,出版杂志的人并不知道谁在看这本杂志,订阅的人并不关心谁在发表这本杂志。出版的人只管将信息发布出去,订阅的人也只在需要的时候收到该信息。
Store-and-forward:存储转发模型类似信件投递,写信的人将消息写给某人,但在将信件发出的时候,收信的人并不一定在家等待,也并不知道有消息给他。但这个消息不会丢失,会放在收信者的信箱中。这种模型允许信息的异步交换。
其他通信模型。。。


Publisher --->Exchange ---> MessageQueue --->Consumer


整个过程是异步的.Publisher,Consumer相互不知道对方的存在,Exchange负责交换/路由,依靠Routing Key,每个消息者有一个Routing Key,每个Binding将自已感兴趣的RoutingKey告诉Exchange,以便Exchange将相关的消息转发给相应的Queue !


几个概念

1
2
3
4
5
6
7
8
9
10
11
12
13
几个概念
Producer,Routing Key,Exchange,Binding,Queue,Consumer.
Producer: 消息的创建者,消息的发送者
Routing Key:唯一用来映射消息该进入哪个队列的标识
Exchange:负责消息的路由,交换
Binding:定义Queue和Exchange的映射关系
Queue:消息队列
Consumer:消息的使用者
Exchange类型
Fan-Out:类似于广播方式,不管RoutingKey
Direct:根据RoutingKey,进行关联投寄
Topic:类似于Direct,但是支持多个Key关联,以组的方式投寄。
       key以.来定义界限。类似于usea.news,usea.weather.这两个消息是一组的。


wKioL1MA096hj-AFAAENzUSW7c0736.jpg


QPID


Qpid 是 Apache 开发的一款面向对象的消息中间件,它是一个 AMQP 的实现,可以和其他符合 AMQP 协议的系统进行通信。Qpid 提供了 C++/Python/Java/C# 等主流编程语言的客户端库,安装使用非常方便。相对于其他的 AMQP 实现,Qpid 社区十分活跃,有望成为标准 AMQP 中间件产品。除了符合 AMQP 基本要求之外,Qpid 提供了很多额外的 HA 特性,非常适于集群环境下的消息通信!


基本功能外提供以下特性:


采用 Corosync(?)来保证集群环境下的Fault-tolerant(?) 特性

支持XML的Exchange,消息为XML时,彩用Xquery过滤

支持plugin

提供安全认证,可对producer/consumer提供身份认证

qpidd --port --no-data-dir --auth

port:端口

--no-data-dir:不指定数据目录

--auth:不启用安全身份认证



启动后自动创建一些Exchange,amp.topic,amp.direct,amp.fanout


tools:


Qpid-config:维护Queue,Exchange,内部配置

Qpid-route:配置broker Federation(联盟?集群?)

Qpid-tool:监控


咱们说完介绍了,这里就赶紧测试下。


服务器端的安装:


1
2
3
yum install qpid-cpp-server
yum install qpid-tools
/etc/init.d/qpidd start


发布端的测试代码

wKioL1MAxo6idPUMAAM2EoX3xTs484.jpg


一些测试代码转自: ibm 的开发社区 


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
#!/usr/bin/env python
#xiaorui.cc
#转自ibm开发社区
import  optparse, time
from qpid.messaging  import  *
from qpid.util  import  URL
from qpid.log  import  enable, DEBUG, WARN
def nameval(st):
   idx = st.find( "=" )
   if  idx >=  0 :
     name = st[ 0 :idx]
     value = st[idx+ 1 :]
   else :
     name = st
     value = None
   return  name, value
parser = optparse.OptionParser(usage= "usage: %prog [options] ADDRESS [ CONTENT ... ]" ,
                                description= "Send messages to the supplied address." )
parser.add_option( "-b" "--broker" default = "localhost" ,
                   help= "connect to specified BROKER (default %default)" )
parser.add_option( "-r" "--reconnect" , action= "store_true" ,
                   help= "enable auto reconnect" )
parser.add_option( "-i" "--reconnect-interval" , type= "float" default = 3 ,
                   help= "interval between reconnect attempts" )
parser.add_option( "-l" "--reconnect-limit" , type= "int" ,
                   help= "maximum number of reconnect attempts" )
parser.add_option( "-c" "--count" , type= "int" default = 1 ,
                   help= "stop after count messages have been sent, zero disables (default %default)" )
parser.add_option( "-t" "--timeout" , type= "float" default =None,
                   help= "exit after the specified time" )
parser.add_option( "-I" "--id" , help= "use the supplied id instead of generating one" )
parser.add_option( "-S" "--subject" , help= "specify a subject" )
parser.add_option( "-R" "--reply-to" , help= "specify reply-to address" )
parser.add_option( "-P" "--property" , dest= "properties" , action= "append" default =[],
                   meta var = "NAME=VALUE" , help= "specify message property" )
parser.add_option( "-M" "--map" , dest= "entries" , action= "append" default =[],
                   meta var = "KEY=VALUE" ,
                   help= "specify map entry for message body" )
parser.add_option( "-v" , dest= "verbose" , action= "store_true" ,
                   help= "enable logging" )
opts, args = parser.parse_args()
if  opts.verbose:
   enable( "qpid" , DEBUG)
else :
   enable( "qpid" , WARN)
if  opts.id  is  None:
   spout_id = str(uuid4())
else :
   spout_id = opts.id
if  args:
   addr = args.pop( 0 )
else :
   parser.error( "address is required" )
content = None
if  args:
   text =  " " .join(args)
else :
   text = None
if  opts.entries:
   content = {}
   if  text:
     content[ "text" ] = text
   for  in  opts.entries:
     name, val = nameval(e)
     content[name] = val
else :
   content = text
conn = Connection(opts.broker,
                   reconnect=opts.reconnect,
                   reconnect_interval=opts.reconnect_interval,
                   reconnect_limit=opts.reconnect_limit)
try :
   conn.open()
   ssn = conn.session()
   snd = ssn.sender(addr)
   count =  0
   start = time.time()
   while  (opts.count ==  0  or count < opts.count) and \
         (opts.timeout  is  None or time.time() - start < opts.timeout):
     msg = Message(subject=opts.subject,
                   reply_to=opts.reply_to,
                   content=content)
     msg.properties[ "spout-id" ] =  "%s:%s"  % (spout_id, count)
     for  in  opts.properties:
       name, val = nameval(p)
       msg.properties[name] = val
     snd.send(msg)
     count +=  1
     print msg
except SendError, e:
   print e
except KeyboardInterrupt:
   pass
conn.close()


客户端的测试代码:

wKiom1MAxznzudhOAARm6bOwuRE241.jpg


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
#!/usr/bin/env python
#xiaorui.cc
##转自ibm开发社区
import  optparse
from  qpid.messaging  import  *
from  qpid.util  import  URL
from  qpid.log  import  enable, DEBUG, WARN
parser  =  optparse.OptionParser(usage = "usage: %prog [options] ADDRESS ..." ,
                                description = "Drain messages from the supplied address." )
parser.add_option( "-b" "--broker" , default = "localhost" ,
                   help = "connect to specified BROKER (default %default)" )
parser.add_option( "-c" "--count" type = "int" ,
                   help = "number of messages to drain" )
parser.add_option( "-f" "--forever" , action = "store_true" ,
                   help = "ignore timeout and wait forever" )
parser.add_option( "-r" "--reconnect" , action = "store_true" ,
                   help = "enable auto reconnect" )
parser.add_option( "-i" "--reconnect-interval" type = "float" , default = 3 ,
                   help = "interval between reconnect attempts" )
parser.add_option( "-l" "--reconnect-limit" type = "int" ,
                   help = "maximum number of reconnect attempts" )
parser.add_option( "-t" "--timeout" type = "float" , default = 0 ,
                   help = "timeout in seconds to wait before exiting (default %default)" )
parser.add_option( "-p" "--print" , dest = "format" , default = "%(M)s" ,
                   help = "format string for printing messages (default %default)" )
parser.add_option( "-v" , dest = "verbose" , action = "store_true" ,
                   help = "enable logging" )
opts, args  =  parser.parse_args()
if  opts.verbose:
   enable( "qpid" , DEBUG)
else :
   enable( "qpid" , WARN)
if  args:
   addr  =  args.pop( 0 )
else :
   parser.error( "address is required" )
if  opts.forever:
   timeout  =  None
else :
   timeout  =  opts.timeout
class  Formatter:
   def  __init__( self , message):
     self .message  =  message
     self .environ  =  { "M" self .message,
                     "P" self .message.properties,
                     "C" self .message.content}
   def  __getitem__( self , st):
     return  eval (st,  self .environ)
conn  =  Connection(opts.broker,
                   reconnect = opts.reconnect,
                   reconnect_interval = opts.reconnect_interval,
                   reconnect_limit = opts.reconnect_limit)
try :
   conn. open ()
   ssn  =  conn.session()
   rcv  =  ssn.receiver(addr)
   count  =  0
   while  not  opts.count  or  count < opts.count:
     try :
       msg  =  rcv.fetch(timeout = timeout)
       print  opts. format  %  Formatter(msg)
       count  + =  1
       ssn.acknowledge()
     except  Empty:
       break
except  ReceiverError, e:
   print  e
except  KeyboardInterrupt:
   pass
conn.close()


Browse 模式的意思是,浏览的意思,一个特殊的需求,我访问了一次,别人也能访问。

Consume 模式的意思是,我浏览了一次后,删除这一条。别人就访问不到啦。

这个是浏览模式:

wKiom1MAx87zwLEFAAbCk_wE0PY041.jpg


sub-pub 通信的例子


Pub-sub 是另一种很有用的通信模型。恐怕它的名字就源于出版发行这种现实中的信息传递方式吧,publisher 就是出版商,subscriber 就是订阅者。


1
2
3
4
5
服务端
qpid-config add exchange topic news-service
./spout news-service/news xiaorui.cc
客户端:
./drain -t  120  news-service/#.news


PUB端,创建TOPIC点 !

wKiom1MAzGDg3nZbAAPRnFen8Y4988.jpg


SUB端,也就是接收端!

wKiom1MAzG2Tsf9YAAO76LbxAS8638.jpg


总结:

 qpid挺好用的,比rabbitmq要轻型,比zeromq保险点! 各方面的文档也都很健全,值得一用。    话说,这三个消息队列我也都用过,最一开始用的是redis的pubsub做日志收集和信息通知,后来在做集群相关的项目的时候,我自己搞了一套zeromq的分布式任务分发,和saltstack很像,当然了远没有万人用的salt强大。  rabbitmq的用法,更是看中他的安全和持久化,当然性能真的不咋地。

关于qpid的性能我没有亲自做大量的测试,但是听朋友说,加持久化可以到7k,不加持久化可以到1500左右。






 本文转自 rfyiamcool 51CTO博客,原文链接:http://blog.51cto.com/rfyiamcool/1359631,如需转载请自行联系原作者


相关文章
|
2月前
|
消息中间件 分布式计算 监控
Python面试:消息队列(RabbitMQ、Kafka)基础知识与应用
【4月更文挑战第18天】本文探讨了Python面试中RabbitMQ与Kafka的常见问题和易错点,包括两者的基础概念、特性对比、Python客户端使用、消息队列应用场景及消息可靠性保证。重点讲解了消息丢失与重复的避免策略,并提供了实战代码示例,帮助读者提升在分布式系统中使用消息队列的能力。
85 2
|
6天前
|
网络协议 Python
python对tcp协议栈进行优化之一
**TCP优化摘要:** - MSS优化涉及调整TCP最大段大小,Python中可使用`socket.getsockopt()`查询MSS。 - Scapy是Python库,用于创建和发送网络包,可用于测试和优化协议栈性能。 - LwIP是轻量级TCP/IP协议栈,适合嵌入式设备,可通过分析和调整提升性能,特别是实时性和资源管理。
|
5天前
|
网络协议 安全 Python
我们将使用Python的内置库`http.server`来创建一个简单的Web服务器。虽然这个示例相对简单,但我们可以围绕它展开许多讨论,包括HTTP协议、网络编程、异常处理、多线程等。
我们将使用Python的内置库`http.server`来创建一个简单的Web服务器。虽然这个示例相对简单,但我们可以围绕它展开许多讨论,包括HTTP协议、网络编程、异常处理、多线程等。
|
5天前
|
Shell 网络安全 数据安全/隐私保护
`paramiko`是一个Python实现的SSHv2协议库
`paramiko`是一个Python实现的SSHv2协议库
|
1月前
|
消息中间件 存储 监控
中间件消息队列协议的缓冲能力
【6月更文挑战第5天】
15 2
|
1月前
|
消息中间件 监控 中间件
|
1月前
|
消息中间件 存储 中间件
中间件消息队列协议异步通信
【6月更文挑战第5天】
21 2
|
1月前
|
消息中间件 存储 中间件
中间件消息队列协议
【6月更文挑战第4天】
20 3
|
2月前
|
网络协议 数据格式 Python
Python进阶---HTTP协议和Web服务器
Python进阶---HTTP协议和Web服务器
30 4
|
1月前
|
消息中间件 存储 RocketMQ
消息队列 MQ产品使用合集之Remoting协议是否可以直接和proxy交互的吗
阿里云消息队列MQ(Message Queue)是一种高可用、高性能的消息中间件服务,它允许您在分布式应用的不同组件之间异步传递消息,从而实现系统解耦、流量削峰填谷以及提高系统的可扩展性和灵活性。以下是使用阿里云消息队列MQ产品的关键点和最佳实践合集。

热门文章

最新文章