【发布时间】:2015-10-04 13:34:15
【问题描述】:
我有这段代码,基本上它运行channel.start_sumption()。 我希望它在一段时间后停止。
我认为 channel.stop_sumption() 是正确的方法:
def stop_consuming(self, consumer_tag=None):
""" Cancels all consumers, signalling the `start_consuming` loop to
exit.
但它不起作用:start_sumption() 永远不会结束(执行不会从此调用退出,“end”永远不会打印)。
导入单元测试 进口鼠兔 导入线程 进口时间
_url = "amqp://user:password@xxx.rabbitserver.com/aaa"
class Consumer_test(unittest.TestCase):
def test_startConsuming(self):
def callback(channel, method, properties, body):
print("callback")
print(body)
def connectionTimeoutCallback():
print("connecionClosedCallback")
def _closeChannel(channel_):
print("_closeChannel")
time.sleep(1)
print("close")
if channel_.is_open:
channel_.stop_consuming()
print("stop_cosuming")
else:
print("channel is closed")
#channel_.close()
params = pika.URLParameters(_url)
params.socket_timeout = 5
connection = pika.BlockingConnection(params)
#connection.add_timeout(2, connectionTimeoutCallback)
channel = connection.channel()
channel.basic_consume(callback,
queue='test',
no_ack=True)
t = threading.Thread(target=_closeChannel, args=[channel])
t.start()
print("start_consuming")
channel.start_consuming() # start consuming (loop never ends)
connection.close()
print("end")
connection.add_timeout 解决我的问题,也许也可以调用 basic_cancel,但我想使用正确的方法。
谢谢
注意: 由于我的声誉点低,我无法对此 (pika, stop_consuming does not work) 做出回应或添加评论。
注2: 我认为我没有跨线程共享通道或连接(Pika 不支持这一点),因为我使用作为参数传递的“通道_”而不是类的“通道”实例(我错了吗?)。
【问题讨论】: