Python Stomp连接心跳超时问题:捕获、避免及on_disconnect未触发排查
Python Stomp心跳超时问题排查与解决
问题背景
使用Python Stomp连接消息中间件时,持续遇到心跳超时问题,接收端后续会陷入空闲状态。当目标端长时间未发送消息后,出现心跳超时错误。
代码框架
def connect_and_subscribe(conn): the_id = 1111 user = 'my_user' password = 'my_password' destination = 'my_destination' subscription_name = 'my_subscriber' conn.connect(login=user,passcode=password, wait=True, wait_time=120, headers = {'client-id': 'xxxx'}) # activemq.subscriptionName ensures durable connection conn.subscribe(destination=destination, id=the_id, ack='auto', persistent=True, headers = {"activemq.subscriptionName":subscription_name, "activemq.persistent":"true"}) class MyListener(stomp.ConnectionListener): def __init__(self, conn, queue): self.conn = conn self.count = 0 self.errors = 0 self.stop = False self.queue = queue def on_error(self, message): print('Received an error %s' % message) self.stop = True def on_message(self, message): if message == "SHUTDOWN": diff = time.time() - self.start print("Received %s in %f seconds" % (self.count, diff)) print("Receiver shutdown") self.stop = True else: item = (self.count, message) self.queue.put(item) self.count += 1 def on_disconnect(self): print_log('Disconnected, going to restart ...') connect_and_subscribe(self.conn) self.stop = False # producer task def producer(queue): host = "my_host" port = my_port destination = "my_destination" heartbeats = 4000 subscription_name = "my_subscriber" conn = stomp.Connection12(host_and_ports = [(host, port)], timeout=120, heartbeats=(heartbeats, heartbeats)) listener = MyListener(conn, queue) conn.set_listener('', listener) print('Listen to ' + host + ':' + str(port) + ' at destination:' + destination) print ('Subscription name:' + subscription_name + ' with heartbeats:' + str(heartbeats)) connect_and_subscribe(conn) while not listener.stop: time.sleep(30) if check_stop(): # a signal will be given if I want to stop the process print('Listener is stopped') listener.stop = True print('Producer is stopped') print('Producer has received ' + str(listener.count) + ' messages in total') queue.put(None) conn.disconnect() print('Producer: Done')
错误信息
heartbeat timeout: diff_receive=6.01600000000326, time=372544.906, lastrec=372538.89
咨询问题
- 如何避免该心跳超时?是否只需增大心跳值?
- 如何捕获该心跳超时,以便进行重连和重订阅?
- 心跳超时发生时,为何
on_disconnect函数未被触发?
解决方案
1. 避免心跳超时的方案
增大心跳值是可选方案之一,但不是唯一解,需结合实际场景调整:
- 调整心跳间隔:当前设置的4000ms(4秒),默认超时阈值为心跳间隔的1.5倍(6秒),无消息时易触发超时。可将心跳值调大,比如
(10000, 10000)(10秒),超时阈值会变为15秒,降低无消息时的超时概率。 - 同步服务端配置:检查消息中间件(如ActiveMQ)的心跳参数,确保服务端心跳发送间隔与客户端匹配,避免服务端未按时发心跳导致客户端误判。
- 优化网络环境:排查是否存在防火墙丢包、网络延迟过高的情况,这类问题也会导致心跳包丢失引发超时。
2. 捕获心跳超时并实现重连
Stomp的ConnectionListener提供on_heartbeat_timeout方法,重写该方法即可捕获心跳超时事件,实现重连逻辑:
修改MyListener类,添加如下方法:
class MyListener(stomp.ConnectionListener): # 原有代码... def on_heartbeat_timeout(self): print_log('心跳超时,开始重连...') # 尝试断开现有连接 try: self.conn.disconnect() except: pass # 重新连接并订阅 connect_and_subscribe(self.conn) self.stop = False
心跳超时发生时,会自动触发该方法执行重连和重订阅操作。
3. on_disconnect未触发的原因
on_disconnect仅在收到服务端发送的DISCONNECT帧时才会触发,而心跳超时是客户端本地检测到的状态,服务端并未主动发送断开指令,因此不会触发该方法。只有服务端主动断开连接、或客户端主动调用conn.disconnect()时,on_disconnect才会被调用。
内容的提问来源于stack exchange,提问作者xymzh
相关产品推荐
相关产品推荐

