使用Stomp.py连接ActiveMQ-STOMP出现消息丢失问题如何排查修复
问题核心原因
- 你对stomp.py的线程模型理解错误:stomp.py的
ConnectionListener的所有回调方法(包括on_message)都是在同一个后台接收线程中串行执行的,不是多线程处理。你在on_message里执行耗时的模型预测逻辑,会直接阻塞整个消息接收线程,导致客户端TCP接收缓冲区堆积溢出,后续消息直接被内核丢弃。 - 你使用的是
Topic(主题)模式而非Queue(队列)模式:ActiveMQ的Topic默认对非持久化订阅的消费者采用"发送即忘"策略,当消费者处理速度跟不上发送速度时,Broker端为该消费者暂存的Pending消息达到内存阈值后,会直接丢弃后续无法投递的消息。 - 你使用了
ack='auto'的确认模式:该模式下消息刚到达客户端、还没完成业务处理就会自动向Broker返回ACK,Broker直接删除该消息,即使后续客户端因为阻塞没来得及处理该消息,也不会有重试机制。 - ActiveMQ STOMP协议默认开启了消费者预取(prefetch)机制,默认预取数量为1000,也就是Broker会一次性给消费者推送1000条消息缓存在客户端,你处理速度过慢时,预取队列满了之后Broker会暂时停止推送,堆积在Broker的Topic消息达到阈值后就会被丢弃。
修复方案
第一步:解耦消息接收和业务处理
不要在on_message里执行耗时的预测逻辑,用本地队列把接收和处理分开,保证接收线程不会被阻塞:
import queue # 初始化本地内存队列,可设置最大长度避免内存溢出 local_queue = queue.Queue(maxsize=10000) def on_message(self, frame): # 接收逻辑只做入队,毫秒级完成,不阻塞接收线程 inputList = frame.body.split(',') local_queue.put((inputList, frame)) # 单独启动工作线程处理模型预测 def worker(): while True: inputList, frame = local_queue.get() try: self.storeInputData(inputList) model.predict(inputList) # 后续补充手动ACK逻辑 except Exception as e: # 异常处理逻辑 pass finally: local_queue.task_done() # 启动多个工作线程(可根据你机器的CPU核心数调整数量) import threading for _ in range(4): t = threading.Thread(target=worker, daemon=True) t.start()
第二步:调整消息确认和预取配置
把ack模式改成client-individual,确保业务处理完成后再手动确认消息,同时降低预取数量避免客户端积压过多消息:
# 订阅时修改配置 self.conn.subscribe( destination='/topic/test', id=1, ack='client-individual', headers={'activemq.prefetchSize': 100} )
在工作线程处理完成后补充手动ACK逻辑:
def worker(): while True: inputList, frame = local_queue.get() try: self.storeInputData(inputList) model.predict(inputList) # 处理成功手动确认 self.conn.ack(frame.headers['message-id'], frame.headers['subscription']) except Exception as e: # 处理失败通知Broker重发或转入死信队列 self.conn.nack(frame.headers['message-id'], frame.headers['subscription']) finally: local_queue.task_done()
第三步:调整消息模式或订阅配置
- 如果你是点对点的消息投递场景,建议直接把
/topic/test换成/queue/test队列模式,队列模式默认会持久化堆积的消息,不会因为消费者处理慢直接丢弃,直到消息被消费或者过期。 - 如果你必须使用Topic模式,建议开启持久化订阅,避免Broker端堆积消息被丢弃:
# 连接时指定唯一客户端ID self.conn.connect('admin', 'admin', wait=True, headers={'client-id': 'unique_consumer_id'}) # 订阅时指定持久化订阅名称 self.conn.subscribe( destination='/topic/test', id=1, ack='client-individual', headers={ 'activemq.prefetchSize': 100, 'activemq.subscriptionName': 'persist_sub_name' } )
第四步:清理发送端冗余配置
发送端不需要订阅目标Topic,删掉发送端的conn.subscribe相关代码即可,避免不必要的资源占用。
内容的提问来源于stack exchange,提问作者dzhou
相关产品推荐
相关产品推荐

