使用Pika停止RabbitMQ Stream队列消费时触发reject/nack不支持错误
问题:RabbitMQ Stream队列中调用stop_consuming导致连接异常
我在应用中使用Stream队列,希望满足特定条件时停止消费。由于使用BlockingConnection时,start_consuming方法会永久阻塞,因此唯一的退出方式是在回调函数中调用stop_consuming。
但这种方式无法生效,原因是根据Pika文档,stop_consuming会拒绝所有待处理消息:
注意:不可确认的待处理消息会丢失;可确认的待处理消息会被拒绝。
而Stream队列不支持拒绝操作,触发如下错误:operation basic.reject caused a connection exception not_implemented: "basic.nack and basic.reject not supported by stream queues queue 'stream-queue' in vhost '/'"
最小复现代例
服务端代码
#!/usr/bin/env python import pika connection = pika.BlockingConnection( pika.ConnectionParameters(host='localhost') ) channel_stream = connection.channel() channel_stream.queue_declare( "stream-queue", auto_delete=False, exclusive=False, durable=True, arguments={ 'x-queue-type': 'stream', } ) channel_stream.basic_qos( prefetch_count=1, ) class Server(object): def __init__(self): channel_stream.basic_consume( queue="stream-queue", on_message_callback=self.stream_callback, ) def stream_callback(self, channel, method, props, body): print(f"received '{body.decode()}' via {method.routing_key}") channel_stream.stop_consuming() server = Server() try: channel_stream.start_consuming() except KeyboardInterrupt: connection.close()
客户端代码
#!/usr/bin/env python import pika connection = pika.BlockingConnection( pika.ConnectionParameters(host='localhost') ) channel_stream = connection.channel() channel_stream.queue_declare( "stream-queue", durable=True, arguments={ 'x-queue-type': 'stream', } ) # 发送两条消息:第一条触发服务端回调,第二条留在队列中,导致服务端被迫nack并崩溃 for i in range(2): channel_stream.basic_publish( exchange='', routing_key='stream-queue', body=f"stream data".encode() ) connection.close()
完整错误输出
~/anaconda3/envs/py310/bin/python ~/workspace/rabbitmq_train/stream_bug/server_stream_only.py received 'stream data' via stream-queue Traceback (most recent call last): File "~/workspace/rabbitmq_train/stream_bug/server_stream_only.py", line 36, in <module> channel_stream.start_consuming() File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 1880, in start_consuming self._process_data_events(time_limit=None) File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 2041, in _process_data_events self.connection.process_data_events(time_limit=time_limit) File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 848, in process_data_events self._dispatch_channel_events() File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 567, in _dispatch_channel_events impl_channel._get_cookie()._dispatch_events() File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 1507, in _dispatch_events consumer_info.on_message_callback(self, evt.method, File "~/workspace/rabbitmq_train/stream_bug/server_stream_only.py", line 30, in stream_callback channel_stream.stop_consuming() File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 1893, in stop_consuming self._cancel_all_consumers() File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 1494, in _cancel_all_consumers self.basic_cancel(consumer_tag) File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 1802, in basic_cancel self._flush_output( File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 1350, in _flush_output self._connection._flush_output(lambda: self.is_closed, *waiters) File "~/anaconda3/envs/py310/lib/python3.10/site-packages/pika/adapters/blocking_connection.py", line 523, in _flush_output raise self._closed_result.value.error pika.exceptions.ConnectionClosedByBroker: (540, "NOT_IMPLEMENTED - basic.nack and basic.reject not supported by stream queues queue 'stream-queue' in vhost '/'") Process finished with exit code 1
启动RabbitMQ命令
docker run -it --rm -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management
内容的提问来源于stack exchange,提问作者tango-taylor
相关产品推荐
相关产品推荐

