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

Crossbar/Autobahn客户端重连后双会话问题技术问询

Crossbar PubSub客户端重连后旧会话残留问题分析

问题场景

我用Crossbar路由实现PubSub功能,基于Autobahn示例开发发布客户端,用Twisted的reactor.callLater实现每5秒定时发布消息。测试自动重连时遇到以下问题:

  • 停止Crossbar路由后,发布操作因传输丢失持续失败
  • 重启路由后,客户端成功重连并生成新会话ID,但旧会话仍在运行,导致新旧两个会话并行发布,旧会话持续失败

我推测问题和reactor.callLater有关,疑惑旧session为何没被覆盖还能运行。虽然可以改用Autobahn的sleep规避,但想了解底层原因。

客户端代码

from autobahn.twisted.component import Component, run
from autobahn.twisted.util import sleep
from twisted.internet.defer import inlineCallbacks
from twisted.internet import reactor
from twisted.internet import defer
from autobahn.wamp.types import PublishOptions
import txaio
import os
import argparse
import six
import treq

txaio.use_twisted()
log = txaio.make_logger()
txaio.start_logging()
#txaio.start_logging(out="jahlog.log", level='info')
url = os.environ.get('CBURL', u'ws://localhost:8080/ws')
realmv = os.environ.get('CBREALM', u'realm1')
topic = os.environ.get('CBTOPIC', 'com.myapp.hello')
topic =  'com.myapp.hello'
print(url, realmv)
component = Component(transports=url, realm=realmv)

@component.on_leave
@inlineCallbacks
def left(session):
    print('session left', session)
    yield session.leave()

@component.on_disconnect
@inlineCallbacks
def gone(session):
    print('session disconnect', session)
    yield session.leave()

@component.on_join
def joined(session, details):
    print("session ready", session)
    sessf = dir(session)
    #print('SESSFFF', sessf)
    startup_system(session)

def finished_pub(res, session):
    print('FINSIED PUB', res)
    #session.leave()
    #session.disconnect()
    reactor.callLater(5, startup_system, session)
     
def failed_pub(res, session):
    print('Failed Pub', res)
    #session.leave()
    reactor.callLater(5, startup_system, session)
 
def startup_system(session):
    """ doc """
    print('START UP SYSTEM', session) 
    d = run_system(session)
    d.addCallback(finished_pub, session)
    d.addErrback(failed_pub, session)
         
@inlineCallbacks
def run_system(session):
    print('START RUN', session)
    sessdir = dir(session)
    #print('SESSDIR', sessdir)
    driver_data = yield prepare_driver_data(session)
    print('driver_data after prepare return') #, driver_data[0][0]) #, driver_data)
    drive_resp = yield publish_driver_data(session, driver_data)
    print('DRIVER PUBLISH SYSTEM GOT DATA', drive_resp)
    return 'OK'

@inlineCallbacks
def prepare_driver_data(session):
    """ doc """
    print('Prepare Drivers data')
    reverse_locations = []
    job_no_saved = None
    driver_url =  b'http://192.168.1.196:8084/loaddrivers'
    headers = {'Content-type': 'application/json'}
    url = driver_url
    driver_data = yield treq.post(url, headers=headers, timeout=30)
    content = yield treq.json_content(driver_data)
    print('DRIVEDR RES: ', content[0][0]) #, content)
    return content

@inlineCallbacks
def publish_driver_data(session, driver_data):
    print('GOT PUB SESS', session)
    print('Do logging', driver_data[0][0])
    options = PublishOptions(acknowledge=True)
    response = yield session.publish(topic, driver_data, options=options)
    print('PUB RESP', response)
    return 'OK'

if __name__ == "__main__":
    #run([component], log_level='info',stop_at_close = True)
    run([component])

底层原因解析

  1. reactor.callLater的引用持有:reactor.callLater会向Twisted事件循环注册延迟任务,该任务会持有传入的session对象引用。只要事件循环未终止,这个引用就会阻止旧session被垃圾回收,任务到期后就会继续执行startup_system,触发旧会话的发布逻辑。

  2. 会话清理不彻底:当会话断开或离开时,你仅调用了session.leave(),但没有取消之前注册的所有reactor.callLater延迟任务。这些旧任务依然存在于事件循环中,会持续触发失效会话的发布流程。

  3. Autobahn重连的会话隔离:Autobahn的Component重连时会创建全新的session实例,但不会自动清理旧会话关联的异步任务。旧会话对象虽已失效,但因延迟任务的引用持有,仍会被执行。

替代解决方案(无需改用Autobahn sleep)

  • 跟踪并取消延迟任务:每次调用reactor.callLater时保存返回的DelayedCall对象,在on_disconnect或on_leave回调中调用其cancel()方法,阻止旧任务继续执行。
  • 维护有效会话引用:用全局变量存储当前有效的会话实例,每次on_join时更新该引用,在发布前检查会话是否为当前有效实例,避免旧会话执行发布。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 01:39:59