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

使用Stomp.py实现AWS ActiveMQ消费者拉取消息并转发到队列

Stomp.py无法拉取队列全量消息解决方案

现有代码的核心问题

  • 订阅未指定ack确认模式:stomp默认使用auto ack,消息到达客户端就自动标记为已消费,若网络波动、消息量较大时易出现未接收就被确认的丢失问题
  • 固定2秒等待时间不合理:队列内消息量较大时,2秒不足以完成全量消息的传输和接收
  • 错误调用了空的listener对象:conn.set_listener()方法无返回值,你定义的listener变量为None,打印listener.message_list属于无效操作
  • 未配置预取参数:AWS MQ(ActiveMQ)默认单消费者预取消息数为1000,若队列消息数远大于该值,会出现分批拉取不及时的问题

修复方案

1. 核心调整点

  • 订阅时指定ack='client-individual',接收并处理完消息后手动调用ack()确认,避免消息丢失
  • 新增动态等待逻辑,替代固定sleep,直到连续指定时间没有新消息接收,再判定为全量拉取完成
  • 订阅时添加预取参数,调整单批拉取的消息数量
  • 修复无效代码逻辑

2. 修正后完整代码

import ssl
import time
import stomp

stompurl = "xxxxxxxx.mq.us-west-2.amazonaws.com"
stompuser = "stomuser"
stomppass = "password"
# 队列名
QUEUE_NAME = "Test_QUEUE"
# 订阅ID(保持唯一即可)
SUBSCRIBE_ID = "102"
# 无新消息超时时间(秒,可根据实际消息量调整)
NO_MSG_TIMEOUT = 5
# 单批预取消息数(根据服务器性能调整,越大拉取越快)
PREFETCH_SIZE = 10000

class MyListener(stomp.ConnectionListener):
    def __init__(self, conn):
        self.msg_list = []
        self.conn = conn
        # 记录最后一次收到消息的时间
        self.last_msg_time = time.time()

    def on_error(self, frame):
        self.msg_list.append(f'(ERROR) {frame.body}')
        self.last_msg_time = time.time()

    def on_message(self, frame):
        self.msg_list.append(frame.body)
        self.last_msg_time = time.time()
        # 处理完消息手动确认,client-individual模式下逐个确认
        self.conn.ack(frame.headers['message-id'], SUBSCRIBE_ID)

# 初始化连接
conn = stomp.Connection(host_and_ports=[(stompurl, 61614)], auto_decode=True)
conn.set_ssl(for_hosts=[(stompurl, 61614)], ssl_version=ssl.PROTOCOL_TLS)
# 传入连接对象给监听器,用于ack确认
lst = MyListener(conn)
conn.set_listener('', lst)
conn.connect(stompuser, stomppass, wait=True)

# 订阅时指定ack模式和预取参数
conn.subscribe(
    destination=QUEUE_NAME,
    id=SUBSCRIBE_ID,
    ack='client-individual',
    headers={'activemq.prefetchSize': PREFETCH_SIZE}
)

# 动态等待:直到连续NO_MSG_TIMEOUT秒没有新消息,判定为全量拉取完成
while True:
    current_time = time.time()
    if current_time - lst.last_msg_time > NO_MSG_TIMEOUT:
        break
    time.sleep(1)

# 拉取到的全量消息
full_messages = lst.msg_list
print(f"共拉取到{len(full_messages)}条消息")
# 此处添加后续将消息推送到目标队列的逻辑
# for msg in full_messages:
#     conn.send(body=msg, destination='YOUR_TARGET_QUEUE')

conn.disconnect()

3. 附加说明

  • 若你需要消费后保留原队列的消息,可将ack模式改为auto,同时去掉手动ack的逻辑,拉取的消息不会标记为已消费,原队列仍会保留消息
  • 若队列消息量极大,建议边拉取边推送,不要全部存在内存中,避免OOM问题
  • 测试时可先在空队列发送少量测试消息,验证接收逻辑正常后再处理全量消息

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 03:06:02