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

ActiveMQ Artemis与stomp.py客户端确认机制不符合预期问题

ActiveMQ Artemis 与 stomp.py 中 ack='client' 模式的 nack 行为异常问题

我在使用ActiveMQ Artemis搭配stomp.py时,发现ack='client'模式的表现不符合预期:调用nack后消息直接消失,并未像预期那样进入DLQ(死信队列);但如果在消费逻辑中主动抛出异常,消息却能正常被转发到DLQ。我的需求是:业务正常时对消息执行ack确认,出现异常时执行nack,让消息进入DLQ。

相关代码

Producer.py

import stomp

# STOMP服务器配置
server = 'localhost'
port = 61613
username = 'artemis'
password = 'artemis'
destination = 'your.queue.name'
message = 'hello :-)'

# 创建STOMP连接
conn = stomp.Connection([(server, port)])
conn.set_listener('', stomp.PrintingListener())

# 连接服务器
conn.connect(username, password, wait=True)

# 发送持久化消息到指定队列
conn.send(destination=destination, body=message, headers={'persistent': 'true'})

# 断开连接
conn.disconnect()

Consumer.py

import stomp
import time

class MyListener(stomp.ConnectionListener):
    def __init__(self, conn):
        self.conn = conn

    def on_error(self, frame):
        print('收到错误:')
        print(f'Headers: {frame.headers}')
        print(f'Body: {frame.body}')

    def on_message(self, frame):
        print('收到消息:')
        print(f'Headers: {frame.headers}')
        print(f'Body: {frame.body}')
        
        time.sleep(1)
        
        # self.conn.ack(frame.headers['message-id'], frame.headers['subscription'])
        
        # 该行代码没有达到预期效果
        self.conn.nack(frame.headers['message-id'], frame.headers['subscription'])
        
        # 抛出异常会将消息送入DLQ
        # raise Exception("我会把这条消息送到DLQ!")
        print('处理完成!!!')

def main():
    server = 'localhost'
    port = 61613
    username = 'artemis'
    password = 'artemis'
    destination = 'your.queue.name'

    conn = stomp.Connection([(server, port)], heartbeats=(4000, 4000))
    conn.set_listener('', MyListener(conn))

    try:
        conn.connect(username, password, wait=True)
    except Exception as e:
        print(f'连接STOMP服务器失败: {e}')
        return

    try:
        # 尝试过的确认模式:
        # client-individual
        # client
        # CLIENT_ACKNOWLEDGE
        conn.subscribe(destination=destination, id='1', ack='client', headers={'subscription-type': 'ANYCAST'})
    except Exception as e:
        print(f'订阅队列失败: {e}')
        conn.disconnect()
        return

    print('等待消息中...')
    
    # 保持主线程存活以接收消息
    try:
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        print('正在断开连接...')
        conn.disconnect()
    except Exception as e:
        print(f'发生意外错误: {e}')
        conn.disconnect()

if __name__ == '__main__':
    main()

broker.xml(相关配置片段)

...
<addresses>
...
    <address name="your.queue.name">
        <anycast>
            <queue name="your.queue.name" />
        </anycast>
        <dead-letter-address>DLQ.your.queue.name</dead-letter-address>
        <expiry-address>ExpiryQueue</expiry-address>
        <redelivery-delay>5000</redelivery-delay> <!-- 重发前延迟5秒 -->
        <max-delivery-attempts>1</max-delivery-attempts> <!-- 最大尝试次数,超过后进入DLQ -->
        <redelivery-delay-multiplier>2.0</redelivery-delay-multiplier> <!-- 指数退避乘数 -->
        <max-redelivery-delay>60000</max-redelivery-delay> <!-- 最大延迟1分钟 -->
    </address>
    
    <address name="DLQ.your.queue.name">
        <anycast>
            <queue name="DLQ.your.queue.name" />
        </anycast>
    </address>
...
</addresses>
...

内容的提问来源于stack exchange,提问作者Demas Sinner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:02:03