我使用以下多线程使用者来解决此问题。
import pika
import time
import sys
import threading
from Queue import Queue
rabbitmq_server = "<RABBITMQ_SERVER_IP>"
queue = "hello1"
connection = None
def check_acknowledge(channel, connection, ack_queue):
delivery_tag = None
while(True):
try:
delivery_tag = ack_queue.get_nowait()
channel.basic_ack(delivery_tag=delivery_tag)
break
except:
connection.process_data_events()
time.sleep(1)
def process_message(body, delivery_tag, ack_queue):
print "Received %s" % (body)
print "Waiting for 600 seconds before receiving next ID\n"
start = time.time()
elapsed = 0
while elapsed < 300:
elapsed = time.time() - start
print "loop cycle time: %f, seconds count: %02d" %(time.clock(), elapsed)
time.sleep(1)
ack_queue.put(delivery_tag)
def callback(ch, method, properties, body):
global connection
ack_queue = Queue()
t = threading.Thread(target=process_message, args=(body, method.delivery_tag, ack_queue))
t.start()
check_acknowledge(ch, connection, ack_queue)
while True:
try:
connection = pika.BlockingConnection(pika.ConnectionParameters(host=rabbitmq_server))
channel = connection.channel()
print ' [*] Waiting for messages. To exit press CTRL+C'
channel.basic_qos(prefetch_count=1)
channel.basic_consume(callback, queue=queue)
channel.start_consuming()
except KeyboardInterrupt:
break
channel.close()
connection.close()
exit(0)
-
消费者
callback
函数触发单独的函数
check_acknowledge
在主线程本身中。因此,连接和通道对象保留在同一线程中。请注意,Pika不是线程安全的,因此我们需要在同一线程中维护这些对象。
-
实际处理发生在从主线程派生的新线程中。
-
一旦
process_message
处理完成后
delivery_tag
在队列中。
-
check\u确认
无限循环,直到找到
delivery\u标签
排队人
process\u消息
。一旦找到,它就会
acks
消息并返回。
我已通过运行此使用者测试了此实现
sleep
持续5分钟、10分钟、30分钟和1小时。这对我来说很有效。