RabbitMQ 生产环境类

更新时间:2020-08-26 10:30:29 点击次数:1231次
一. RabbitMq公共方法类
import datetime
 
class RabbitMQ(object):
    '''rabbitmq 公共配置'''
 
    def __init__(self, queue_name='', exchange='', exchange_type='', routing_key=''):
        self.connection = self.get_connection()
        # 队列名字
        self.queue_name = queue_name  if queue_name else 'handle_info_{}'.format(datetime.datetime.now().strftime("%Y%m%d%H:%M:%S"))
        # 交换机
        self.exchange = exchange if exchange else self.queue_name
        # exchange 类型
        self.exchange_type = exchange_type if exchange_type else 'direct'
        # routing_key
        self.routing_key = routing_key if routing_key else self.queue_name
 
 
    def get_connection(self):
        '''连接方法'''
        host = self.RABBITMQ_HOST = '127.0.0.1'
        port = self.RABBITMQ_POST = 5673
        user = self.RABBITMQ_USER = 'guest'
        password = self.RABBITMQ_PASSWORD = 'guest'
        connection = pika.BlockingConnection(pika.ConnectionParameters(host=host,port=port,credentials=pika.PlainCredentials(user,password),virtual_host="oms"))
        return connection
 
 
    def send(self, body):
        """
        :param body: 消息体
        :param delay: 是否延迟发送
        :param ttl: 延迟时间
        """
        # 声明收容队列
        channel_state = self.connection.channel()
        channel_state.queue_declare(queue=self.queue_name, durable=True)  # durable 消息持久化
        # 声明收容交换机
        channel_state.exchange_declare(exchange=self.exchange, exchange_type=self.exchange_type)
        # 收容队列和收容交换机绑定
        channel_state.queue_bind(exchange=self.exchange, queue=self.queue_name, routing_key=self.routing_key)
 
        for item in body:
            message = json.dumps(item).encode()
            channel_state.basic_publish(exchange=self.exchange, routing_key=self.routing_key, body=message,
                                   properties=pika.BasicProperties(delivery_mode=2))   # make message persistent
 
        self.connection.close()
二.消费核心代码
def callback(ch, method, properties, body):
    '''消息处理体 '''
    print " [x] Received %r" % (body,)
    print " [x] Done"
    ch.basic_ack(delivery_tag = method.delivery_tag)   # 消息确认(异常中断会重新发给消费者消费)
 
channel.basic_qos(prefetch_count=1)   # 公平调度(处理能力强的消费者处理更多的消息)
channel.basic_consume(queue_name,callback,True)
 

本站文章版权归原作者及原出处所有 。内容为作者个人观点, 并不代表本站赞同其观点和对其真实性负责,本站只提供参考并不构成任何投资及应用建议。本站是一个个人学习交流的平台,网站上部分文章为转载,并不用于任何商业目的,我们已经尽可能的对作者和来源进行了通告,但是能力有限或疏忽,造成漏登,请及时联系我们,我们将根据著作权人的要求,立即更正或者删除有关内容。本站拥有对此声明的最终解释权。

回到顶部
嘿,我来帮您!