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])
底层原因解析
reactor.callLater的引用持有:reactor.callLater会向Twisted事件循环注册延迟任务,该任务会持有传入的session对象引用。只要事件循环未终止,这个引用就会阻止旧session被垃圾回收,任务到期后就会继续执行startup_system,触发旧会话的发布逻辑。会话清理不彻底:当会话断开或离开时,你仅调用了
session.leave(),但没有取消之前注册的所有reactor.callLater延迟任务。这些旧任务依然存在于事件循环中,会持续触发失效会话的发布流程。Autobahn重连的会话隔离:Autobahn的
Component重连时会创建全新的session实例,但不会自动清理旧会话关联的异步任务。旧会话对象虽已失效,但因延迟任务的引用持有,仍会被执行。
替代解决方案(无需改用Autobahn sleep)
- 跟踪并取消延迟任务:每次调用
reactor.callLater时保存返回的DelayedCall对象,在on_disconnect或on_leave回调中调用其cancel()方法,阻止旧任务继续执行。 - 维护有效会话引用:用全局变量存储当前有效的会话实例,每次
on_join时更新该引用,在发布前检查会话是否为当前有效实例,避免旧会话执行发布。
内容的提问来源于stack exchange,提问作者John Aherne
相关产品推荐
相关产品推荐

