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

Flet技术问询:如何在应用外部发送PubSub广播消息

问题描述

我正在构建一款Flet应用,用于控制硬件设备并显示从设备接收的数据。

Block Diagram

除桌面应用外,我还希望支持多用户通过多会话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()

关键修改说明

  1. 硬件接口类新增会话管理:

    • 添加active_pages列表存储所有活跃会话的page对象;
    • 新增page_lock保证线程安全地修改/遍历列表;
    • 实现add_page和remove_page方法管理会话生命周期。
  2. 定时线程广播逻辑:

    • 在cb_timer中,先复制active_pages避免遍历过程中列表被修改;
    • 遍历每个page调用pubsub.send_all()发送数据;
    • 捕获异常并移除无效的page(比如会话已关闭)。
  3. 会话生命周期绑定:

    • 在session_main中,会话创建时调用g_hw_if.add_page(page);
    • 设置page.on_disconnect回调,会话关闭时移除page。
  4. 兼容性优化:

    • 移除view=ft.WEB_BROWSER,让Flet自动根据运行环境选择桌面或Web视图,同时保留port参数方便Web访问。

这样修改后,硬件接口的定时线程只需发起一次数据请求,然后向所有活跃会话广播数据,同时兼容桌面和Web多会话场景,避免了会话单独请求的冗余问题,也解决了首个会话关闭后硬件通信中断的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 09:07:05