rabbitMQ玩耍日记

2019-10-25  本文已影响0人  魂狩

rabbitMQ是个啥东东,想要玩的人自然已经有了初步了解。这里直接开始安装。安装环境为ubuntu 18.04。

  1. ubuntu自带rabbitMQ的包,直接apt安装即可。apt install rabbitmq-server。安装时会把rabbitmq和其依赖erlang一起安装。
  2. 默认情况下web管理界面是不会开启的,需要手动启动。rabbitmq-plugins enable rabbitmq_management
  3. 启动后默认端口为15672,默认用户guest/guest。登录进去后可以先创建自己的帐号,然后将默认帐号删除;也可以直接使用命令行进行操作:
    1. rabbitmqctl add_user small 123456 创建一个密码是123456的small用户
    2. rabbitmqctl set_user_tags small administrator 给此用户授管理员权限。
  4. 创建用户后需要给用户权限。根据情况给用户分配虚拟主机(virtual hosts),用户只能在对应虚拟主机进行操作。
  5. 分配权限后就可以愉快玩耍了。以下为最简单的demo,来自官方文档:
#生产者
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters(credentials=pika.PlainCredentials('small','123456')))
channel = connection.channel()
channel.queue_declare(queue='test')
channel.basic_publish(exchange='', routing_key='test', body='hello world')
connection.close()
#消费者
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters(credentials=pika.PlainCredentials('small','123456')))
channel = connection.channel()
channel.queue_declare(queue='test')
def callback(ch, method, properties, body):
    print(" [x] Received %r" % body)

channel.basic_consume(
    queue='test', on_message_callback=callback, auto_ack=True)

print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()

其他说明

多线程处理

pika是非多线程安全的!!!
所以如果想多线程处理消息,要么每个线程建立一个连接(Connections),要么使用python自带的queue来做缓冲。每个线程建立一个连接开销比较大,但是代码简单,在压力不大的情况下足够用了。使用内置队列,代码稍微复杂一点,但是性能更好。以下为内置队列的demo。

import pika,functools,queue
connection = pika.BlockingConnection(pika.ConnectionParameters(credentials=pika.PlainCredentials('small','123456')))
channel = connection.channel()
channel.queue_declare(queue='test')
class worker():
    q = queue.Queue(10)

    @staticmethod
    def work():
        while True:
            data = worker.q.get()
            body=data['body']
            # do something
            cutil.rabbit.add_callback_threadsafe(functools.partial(channel.basic_ack, data['tag']))

def rec(ch, method, properties, body):
    try:
        body = json.loads(body)
        worker.q.put({'body': body, 'tag': method.delivery_tag})
    except Exception as e:
        print('sender ???',end='')
        print(e,flush=True)

for x in range(10):
    threading.Thread(target=worker.work, daemon=True).start()
init()
channel.basic_consume(queue=RABBITC.SEND.value, on_message_callback=rec)
channel.start_consuming()

注意:在子线程不能直接调用channel.basic_ack,必须使用add_callback_threadsafe间接执行。更普遍的,在子线程不能直接调用channel的其他任何函数,都需要间接执行,add_callback_threadsafe会让channel在主线程执行对应函数。

心跳机制

在连接的时候,只要ConnectionParameters中的heartbeat不为0,rabbitMQ就会启用心跳机制,在超时没有数据且没有收到心跳包的情况下断开连接。
作为客户端(调用了start_consuming方法),pika在没有接收数据的时候,间隔一段时间会主动发送心跳包,维持连接。然而,作为服务端,pika没有事件时是不会发送心跳包的。这样会引起一个问题:当业务量很少,比如间隔几分钟才添加一条消息,pika却不会主动发心跳包。在中间等待数据的间隔,rabbitMQ会超时断开连接,pika再尝试添加消息时会出错。
解决方案有两种,一种是重连机制,发送失败就重新连接然后再发送。另外一种是手动发送心跳包,自己设置一个定时器,到时间了执行channel.queue_declare命令,pika就会正常发心跳包了。

上一篇 下一篇

猜你喜欢

热点阅读