使用Stomp.py实现AWS ActiveMQ消费者拉取消息并转发到队列
Stomp.py无法拉取队列全量消息解决方案
现有代码的核心问题
- 订阅未指定
ack确认模式:stomp默认使用autoack,消息到达客户端就自动标记为已消费,若网络波动、消息量较大时易出现未接收就被确认的丢失问题 - 固定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__
相关产品推荐
相关产品推荐

