Flet技术问询:如何在应用外部发送PubSub广播消息
我正在构建一款Flet应用,用于控制硬件设备并显示从设备接收的数据。

除桌面应用外,我还希望支持多用户通过多会话Web应用访问。我在全局作用域中定义了一个硬件接口类实例:会话通过带线程安全锁的write函数向硬件发送指令(蓝色箭头);目前从硬件接收数据的方式是会话通过点击事件调用带锁的read函数,请求数据后将结果返回给对应会话。
现在我需要实现硬件数据的周期性请求:不希望每个会话单独发起请求(如3个会话会发送3次请求、接收3次回复),而是在硬件接口类中设置定时线程,仅发送1次请求,接收回复后通过PubSub向所有会话广播数据(红色箭头)。
请问是否可以在Flet应用外部(定时线程中)使用PubSub向所有会话发送广播消息?或者是否可以创建一个无界面的无头会话,在其中嵌入硬件接口类,使用page.pubsub.send_all()完成广播?希望代码同时兼容桌面应用。
我曾尝试使用Flet包中的PubSub实例,但发现它需要WebSocket连接才能与会话通信,不知如何配置;也曾考虑让首个会话作为硬件通信的专属实例,但该会话关闭后其他会话将无法与硬件通信。
极简示例代码
import flet as ft from datetime import datetime import threading class PeriodicThread(object): def __init__(self, callback=None, period=1, name=None, *args, **kwargs): self.name = name self.args = args self.kwargs = kwargs self.callback = callback self.period = period self.stop = False self.current_timer = None self.schedule_lock = threading.Lock() def set_period(self, period): self.period = period def start(self): self.stop = False self.schedule_timer() def run(self): if self.callback is not None: self.callback(*self.args, **self.kwargs) def _run(self): try: self.run() except Exception as e: pass finally: with self.schedule_lock: if not self.stop: self.schedule_timer() def schedule_timer(self): self.current_timer = threading.Timer(self.period, self._run, *self.args, **self.kwargs) if self.name: self.current_timer.name = self.name self.current_timer.start() def cancel(self): with self.schedule_lock: self.stop = True if self.current_timer is not None: self.current_timer.cancel() def join(self): self.current_timer.join() class hardware_interface(): def __init__(self): self.lock = threading.Lock() self.cnt = 0 self.timer = PeriodicThread(callback=self.cb_timer, period=1, name="timer_periodic_update") self.timer.start() def send_command(self, cmd, cmd_data): print(f'Instrument: executing command "{cmd} {cmd_data}"') def read_oneoff_data(self): data = self.cnt self.cnt += 1 print(f'Instrument: data requested, replied {data}') return data def cb_timer(self): data = datetime.now() print(f'Instrument: data requested, replied {data}') #TODO: Send broadcast message with the data to all app sessions # If this would be a flet page object I could use # self.pubsub.send_all(data) class App(ft.UserControl): def build(self): self.tf_cmd = ft.TextField(label="CMD") self.tf_cmd_data = ft.TextField(label="CMD Data") self.tf_req_data = ft.TextField(label="Requested Data") self.tf_per_data = ft.TextField(label="Periodic Data") page = ft.Column( [ ft.Row( [ self.tf_cmd, self.tf_cmd_data, ft.ElevatedButton(text="Send Command", on_click=self.btn_send_click) ] ), ft.Row( [ self.tf_req_data, ft.ElevatedButton(text="Request Data", on_click=self.btn_reqdata_click) ] ), ft.Row( [ self.tf_per_data ] ) ] ) self.page.pubsub.subscribe(self.on_HW_periodic_data_msg) return page def btn_send_click(self, e): g_hw_if.lock.acquire(timeout=3) g_hw_if.send_command(self.tf_cmd.value, self.tf_cmd_data.value) g_hw_if.lock.release() def btn_reqdata_click(self, e): g_hw_if.lock.acquire(timeout=3) self.tf_req_data.value = g_hw_if.read_oneoff_data() g_hw_if.lock.release() self.update() def on_HW_periodic_data_msg(self, msg): self.tf_per_data.value = msg self.update() def session_main(page: ft.Page): app = App() page.add(app) def main(): global g_hw_if g_hw_if = hardware_interface() ft.app(target=session_main, view=ft.WEB_BROWSER, port=45678) if __name__ == '__main__': main()
运行代码后,浏览器会打开Web应用标签页,打开新标签页访问同一URL即可启动新会话。可在任意会话中发送指令或请求数据,hardware_interface类的cb_timer方法会周期性获取数据,需在#TODO处实现向所有Flet会话广播数据的功能。
要实现硬件接口定时线程向所有Flet会话广播数据,同时兼容桌面和Web应用,可通过以下方式修改代码:
核心思路
Flet的PubSub依赖会话的WebSocket连接,无法直接在外部线程调用。我们可以:
- 在全局维护一个所有活跃会话的
page对象列表,通过线程安全的方式管理这个列表; - 硬件接口的定时线程遍历列表,对每个
page调用pubsub.send_all()发送数据; - 在会话创建/销毁时自动添加/移除对应的
page对象,避免无效连接。
修改后的完整代码
import flet as ft from datetime import datetime import threading class PeriodicThread(object): def __init__(self, callback=None, period=1, name=None, *args, **kwargs): self.name = name self.args = args self.kwargs = kwargs self.callback = callback self.period = period self.stop = False self.current_timer = None self.schedule_lock = threading.Lock() def set_period(self, period): self.period = period def start(self): self.stop = False self.schedule_timer() def run(self): if self.callback is not None: self.callback(*self.args, **self.kwargs) def _run(self): try: self.run() except Exception as e: pass finally: with self.schedule_lock: if not self.stop: self.schedule_timer() def schedule_timer(self): self.current_timer = threading.Timer(self.period, self._run, *self.args, **self.kwargs) if self.name: self.current_timer.name = self.name self.current_timer.start() def cancel(self): with self.schedule_lock: self.stop = True if self.current_timer is not None: self.current_timer.cancel() def join(self): self.current_timer.join() class hardware_interface(): def __init__(self): self.lock = threading.Lock() self.cnt = 0 # 全局会话page列表 + 线程锁 self.active_pages = [] self.page_lock = threading.Lock() self.timer = PeriodicThread(callback=self.cb_timer, period=1, name="timer_periodic_update") self.timer.start() def send_command(self, cmd, cmd_data): print(f'Instrument: executing command "{cmd} {cmd_data}"') def read_oneoff_data(self): data = self.cnt self.cnt += 1 print(f'Instrument: data requested, replied {data}') return data def cb_timer(self): data = datetime.now().strftime("%Y-%m-%d %H:%M:%S") print(f'Instrument: data requested, replied {data}') # 线程安全地遍历所有活跃page,广播数据 with self.page_lock: # 复制列表避免遍历过程中列表被修改 pages_copy = self.active_pages.copy() for page in pages_copy: try: # 直接向page的pubsub发送消息 page.pubsub.send_all(data) except Exception as e: # 移除无效的page连接 with self.page_lock: if page in self.active_pages: self.active_pages.remove(page) # 添加会话page def add_page(self, page): with self.page_lock: if page not in self.active_pages: self.active_pages.append(page) # 移除会话page def remove_page(self, page): with self.page_lock: if page in self.active_pages: self.active_pages.remove(page) class App(ft.UserControl): def build(self): self.tf_cmd = ft.TextField(label="CMD") self.tf_cmd_data = ft.TextField(label="CMD Data") self.tf_req_data = ft.TextField(label="Requested Data") self.tf_per_data = ft.TextField(label="Periodic Data") page = ft.Column( [ ft.Row( [ self.tf_cmd, self.tf_cmd_data, ft.ElevatedButton(text="Send Command", on_click=self.btn_send_click) ] ), ft.Row( [ self.tf_req_data, ft.ElevatedButton(text="Request Data", on_click=self.btn_reqdata_click) ] ), ft.Row( [ self.tf_per_data ] ) ] ) self.page.pubsub.subscribe(self.on_HW_periodic_data_msg) return page def btn_send_click(self, e): g_hw_if.lock.acquire(timeout=3) g_hw_if.send_command(self.tf_cmd.value, self.tf_cmd_data.value) g_hw_if.lock.release() def btn_reqdata_click(self, e): g_hw_if.lock.acquire(timeout=3) self.tf_req_data.value = g_hw_if.read_oneoff_data() g_hw_if.lock.release() self.update() def on_HW_periodic_data_msg(self, msg): self.tf_per_data.value = msg self.update() def session_main(page: ft.Page): # 会话创建时添加page到硬件接口的活跃列表 g_hw_if.add_page(page) app = App() page.add(app) # 会话销毁时移除page def on_page_disconnect(e): g_hw_if.remove_page(page) page.on_disconnect = on_page_disconnect def main(): global g_hw_if g_hw_if = hardware_interface() # 兼容桌面和Web,根据运行环境自动切换view ft.app(target=session_main, port=45678) if __name__ == '__main__': main()
关键修改说明
硬件接口类新增会话管理:
- 添加
active_pages列表存储所有活跃会话的page对象; - 新增
page_lock保证线程安全地修改/遍历列表; - 实现
add_page和remove_page方法管理会话生命周期。
- 添加
定时线程广播逻辑:
- 在
cb_timer中,先复制active_pages避免遍历过程中列表被修改; - 遍历每个
page调用pubsub.send_all()发送数据; - 捕获异常并移除无效的
page(比如会话已关闭)。
- 在
会话生命周期绑定:
- 在
session_main中,会话创建时调用g_hw_if.add_page(page); - 设置
page.on_disconnect回调,会话关闭时移除page。
- 在
兼容性优化:
- 移除
view=ft.WEB_BROWSER,让Flet自动根据运行环境选择桌面或Web视图,同时保留port参数方便Web访问。
- 移除
这样修改后,硬件接口的定时线程只需发起一次数据请求,然后向所有活跃会话广播数据,同时兼容桌面和Web多会话场景,避免了会话单独请求的冗余问题,也解决了首个会话关闭后硬件通信中断的问题。
内容的提问来源于stack exchange,提问作者midiwidi

