温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

python rabbitmq 消费端根据能力轮询接受

发布时间:2020-07-25 16:25:44 来源:网络 阅读:862 作者:w411639146 栏目:建站服务器

给接收端添加:

channel.basic_qos(prefetch_count=1)  ##一次处理一个,处理完再接受新消息



发送端:

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()


channel.queue_declare(queue='hello',durable=True)  ##队列持久化,队列重启后也存在,不保证数据是否存在
# channel.queue_delete(queue="task_queue")
for i in range(100):
    channel.basic_publish(exchange='',
                          routing_key='hello',
                          body=str(i),
                          properties=pika.BasicProperties(delivery_mode=2) ##数据持久化
                          )
# print("Sent 'hello world!'")
connection.close()


接收端:

#!/usr/bin/env python
import pika
import time
connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
channel = connection.channel()

channel.queue_declare(queue='hello',durable=True)
# channel.queue_bind(queue='hello',exchange='',routing_key='hello')
def callback(ch, method, properties, body):
    # print("aaa")
    print(" [x] Received %r" % body)
    time.sleep(1)
    ch.basic_ack(delivery_tag=method.delivery_tag)  # 给rabbitmq返回已拿到数据信号。

channel.basic_qos(prefetch_count=1)  ##一次处理一个,处理完再接受新消息
channel.basic_consume(callback,
                      queue='hello',
                      no_ack=False)

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


向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI