RabbitMQ的负载均衡(Python)
参考:https://www.rabbitmq.com/tutorials/tutorial-two-python.html
默认分发方式
先观察一下当存在多个消息接收者时,RabbitMQ的处理方式。
基于前面的Hello World代码,修改send.py,另存为new_tasks.py:
new_tasks.py
修改receive.py,另存为worker.py:
worker.py
打开两个终端窗口,运行python worker.py:
再新建一个窗口,用于发送消息:
接收端的结果为:
可见RabbitMQ默认将消息均匀分发到两个worker(截图的时候忘了测试带"..."模拟耗时任务,结果是一样的)。
应用系统中,如果奇数消息耗时较多,且全部分发到窗口1,偶数消息耗时较少,全部分发到窗口2,这显然是不合理的,需要让程序拥有自动分配任务的功能。
消息确认
消息确认用于消息接收者,在接收到消息后发回一个接收到消息的确认消息。这一点用于后面的操作。
修改worker.py文件,去掉no_ack=true参数,并在回调函数中加入一条语句:
有了消息确认后,如果程序终止,消息会被重发到另一个接收者,比如上面两个接收者窗口的情况,一个包含1、3,另一个包含2、4,关闭第一个窗口中的程序,第二个窗口的情况:
再次运行一个新的接收者程序,终止此时包含2、4、1、3的窗口程序,新窗口会接收到所有消息:
消息自动按照时间排序。
当然这还不能解决负载均衡的问题。
数据持久化
现在的消息记录并没有被系统保存,终止程序,消息被销毁。只需要在声明消息队列时加上durable=Ture参数,消息便会被RabbitMQ系统保存。
这样写程序会报错,因为我们之前定义过不带持久化参数的hello队列,系统不允许重复声明队列。可以更改队列名称来解决这个问题,注意new_tasks和worker都要更改:
这只是声明了队列属性,完成消息队列持久化还需要在发送消息时传入参数,通知系统这条消息需要持久化:
到此也仍然没有解决负载均衡的问题。
负载均衡
我们就是想让系统分发消息的时候智能一点,分配消息任务给较闲置的处理者,那让接收者主动告诉系统它是否有空,问题不就好解决了吗。
将这条语句放在basic_consum前面。这个函数设置了prefetch_count参数为1,是告诉RabbitMQ它一次只处理一条消息,并且系统在接收到响应之前不要给它分发消息。这就是前面消息确认的作用。
现在RabbitMQ可以公平地分发消息队列了。