You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 10:01:08