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

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

咨询问题

  1. 如何避免该心跳超时?是否只需增大心跳值?
  2. 如何捕获该心跳超时,以便进行重连和重订阅?
  3. 心跳超时发生时,为何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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 16:25:37