Textual 0.18.0 并发管理 Worker API:统一管理 asyncio 任务与线程
【免费下载链接】textualThe lean application framework for Python. Build sophisticated user interfaces with a simple Python API. Run your apps in the terminal and a web browser.项目地址: https://gitcode.com/gh_mirrors/te/textual
导读
Textual 0.18.0 于 2023 年 4 月正式引入全新的 Worker API,为构建终端与浏览器界面的应用提供了一套统一的并发管理方案——无论你的后台任务是 asyncio 协程还是普通函数线程,都可以用同一个run_worker方法与@work装饰器来启动、跟踪、取消并有序关闭。读完本文,你将掌握如何把耗时的网络请求、子进程调用或计算密集任务移出消息循环,让界面保持即时响应,并学会用状态机、事件与线程安全机制稳健地组织应用并发逻辑。
版本背景:为什么 UI 框架需要自己的并发管理 API
Textual 0.18.0 距离上一版本发布不足一周,核心变化只有一个,却直击大量真实应用的痛点。正如 发布说明 所述:人们正在用 Textual 构建对接 REST API、WebSocket 和外部进程的应用,而他们遇到的是 asyncio 任务与线程的通用并发难题——界面卡顿、响应乱序、任务无法有序关闭。Textual 团队认为仅把用户指向 asyncio 文档是不够的,需要一个更好的答案,于是诞生了 Worker API。
一个容易被忽略的细节是:Textual 此前已经实现了一套有序关闭机制来清理支撑 Widget 的任务——子节点先于父节点关闭,逐级上溯直到 App(根节点)。新的 Worker API 直接复用了这套机制,保证由 Widget、Screen 或 App 发起的 Worker 任务以同样的顺序被关闭,从而规避了异步代码中最棘手的“任务泄漏”与“退出竞态”问题。
核心新增:run_worker 与 Worker 对象
run_worker是本次发布的入口方法,它定义在 src/textual/dom.py 中,因此任何 DOM 节点(Widget、Screen、App)都可以直接调用。它接收一个函数、协程或 awaitable,把它包装成 Worker 对象在后台运行,并返回该对象供你跟踪。
让我们用一个天气查询应用来直观理解它解决的问题。weather01.py是“反例”:在on_input_changed消息处理器里直接await网络请求:
import httpx from rich.text import Text from textual.app import App, ComposeResult from textual.containers import VerticalScroll from textual.widgets import Input, Static class WeatherApp(App): """App to display the current weather.""" CSS_PATH = "weather.tcss" def compose(self) -> ComposeResult: yield Input(placeholder="Enter a City") with VerticalScroll(id="weather-container"): yield Static(id="weather") async def on_input_changed(self, message: Input.Changed) -> None: """Called when the input changes""" await self.update_weather(message.value) async def update_weather(self, city: str) -> None: """Update the weather for the given city.""" weather_widget = self.query_one("#weather", Static) if city: # Query the network API url = f"https://wttr.in/{city}" async with httpx.AsyncClient() as client: response = await client.get(url) weather = Text.from_ansi(response.text) weather_widget.update(weather) else: # No city, so just blank out the weather weather_widget.update("") if __name__ == "__main__": app = WeatherApp() app.run()运行这个版本时你会发现:输入框明显不跟手,按键到回显之间有一两百毫秒甚至数秒的延迟——因为整个应用在请求完成前无法处理其他任何消息。修复方式只有一行改动(weather02.py第 21 行):
import httpx from rich.text import Text from textual.app import App, ComposeResult from textual.containers import VerticalScroll from textual.widgets import Input, Static class WeatherApp(App): """App to display the current weather.""" CSS_PATH = "weather.tcss" def compose(self) -> ComposeResult: yield Input(placeholder="Enter a City") with VerticalScroll(id="weather-container"): yield Static(id="weather") async def on_input_changed(self, message: Input.Changed) -> None: """Called when the input changes""" self.run_worker(self.update_weather(message.value), exclusive=True) async def update_weather(self, city: str) -> None: """Update the weather for the given city.""" weather_widget = self.query_one("#weather", Static) if city: # Query the network API url = f"https://wttr.in/{city}" async with httpx.AsyncClient() as client: response = await client.get(url) weather = Text.from_ansi(response.text) weather_widget.update(weather) else: # No city, so just blank out the weather weather_widget.update("") if __name__ == "__main__": app = WeatherApp() app.run()run_worker会立即调度update_weather并返回 Worker 对象,消息循环不会被阻塞;协程随后在后台并发运行,一两秒后完成。注意这里还传了exclusive=True,它解决了一个隐蔽问题:并发请求的响应到达顺序不保证与请求顺序一致。比如你输入 "Paris",可能 "Pari" 的响应比 "Paris" 更晚到达,导致显示错误城市的数据。exclusive标志会让 Textual 在新 Worker 启动前取消所有同组的旧 Worker(见 src/textual/worker_manager.py 中cancel_group的调用逻辑),保证只有最后一次输入的结果被展示。
run_worker 与 @work 的完整参数
从 dom.py 的 run_worker 定义 和 _work_decorator.py 可以看到,两个入口接受几乎一致的参数(@work多一个由run_worker内部处理、run_worker多一个start):
| 参数 | 默认值 | 作用 |
|---|---|---|
name | ""(取自方法名) | 短字符串,用于日志与调试中标识 Worker |
group | "default" | 短字符串,标识一组 Worker;exclusive=True时按 group 取消 |
exit_on_error | True | 出错时是否直接退出应用并打印 traceback;设为False可抑制异常 |
exclusive | False | 启动前取消同组所有 Worker |
description | 方法名+实参 | 更长的可读描述,便于调试;超过 1000 字符会被截断 |
thread | False | 标记为线程 Worker(普通函数必须开启) |
start(仅run_worker) | True | 是否立即启动 Worker |
@work 装饰器:把协程与普通函数统一成 Worker
run_worker需要你手动调用;而 @work 装饰器 更进一步——直接改造被装饰的方法本身,调用该方法时自动创建并启动 Worker。发布说明特别强调:它能同时接受协程和普通函数,分别以 asyncio 任务或线程的方式调度。
weather03.py展示了它的用法:
import httpx from rich.text import Text from textual import work from textual.app import App, ComposeResult from textual.containers import VerticalScroll from textual.widgets import Input, Static class WeatherApp(App): """App to display the current weather.""" CSS_PATH = "weather.tcss" def compose(self) -> ComposeResult: yield Input(placeholder="Enter a City") with VerticalScroll(id="weather-container"): yield Static(id="weather") async def on_input_changed(self, message: Input.Changed) -> None: """Called when the input changes""" self.update_weather(message.value) @work(exclusive=True) async def update_weather(self, city: str) -> None: """Update the weather for the given city.""" weather_widget = self.query_one("#weather", Static) if city: # Query the network API url = f"https://wttr.in/{city}" async with httpx.AsyncClient() as client: response = await client.get(url) weather = Text.from_ansi(response.text) weather_widget.update(weather) else: # No city, so just blank out the weather weather_widget.update("") if __name__ == "__main__": app = WeatherApp() app.run()添加@work(exclusive=True)后,update_weather的调用点不再需要await(见第 22 行)——尽管它仍是async def。装饰器接收与run_worker相同的参数,并在内部转调self.run_worker(partial(method, *args, **kwargs), ...)(见 _work_decorator.py),因此被装饰方法必须定义在 DOMNode 子类上。
需要特别留意一个约束:源码中有明确校验——
Textual 会对“普通函数 + 未设置
thread=True”的组合抛出WorkerDeclarationError(“Can not create a worker from a non-async function unlessthread=Trueis set”),见 src/textual/_work_decorator.py。
即:非 async 函数必须声明为线程 Worker,否则启动即报错。
读取 Worker 的返回值
Worker 函数的返回值只有在任务完成后才可用,通过worker.result属性访问——初始为None,完成后替换为真实返回值(见 src/textual/worker.py)。
如果你确实需要同步等待结果,可以调用worker.wait()协程(src/textual/worker.py)。但要注意:在消息处理器里await worker.wait()同样会阻塞界面更新,与最初的问题别无二致。因此更推荐的方式是监听 Worker 事件 或定期轮询worker.state。另外,wait()内部有防死锁保护:若在 Worker 自身函数内调用wait(),会抛出DeadlockError(src/textual/worker.py)。
取消 Worker
在任务完成前,随时可以调用worker.cancel()取消它(src/textual/worker.py)。对协程 Worker,这会在协程内部抛出asyncio.CancelledError,使其提前退出;同时会设置cancelled_event(一个 threading.Event),供线程 Worker 轮询判断。注意is_cancelled属性(src/textual/worker.py)表示“已发出取消请求”,被取消的任务可能仍在运行中。
异常处理:exit_on_error 的两种行为
Worker 抛出异常时,默认行为是直接退出应用并在终端打印 traceback——这对调试期很友好,但生产环境可能需要更优雅的降级。通过run_worker(..., exit_on_error=False)或@work(exit_on_error=False)可以创建“不因异常退出”的 Worker,此时异常对象会保存在worker.error中(src/textual/worker.py),你可以在事件处理器里自行决定如何处理。
Worker 生命周期:归属节点与自动清理
Worker 由应用内唯一的 WorkerManager 实例统一管理,可通过app.workers(或任一 DOM 节点的workers属性)访问。它是一个容器类对象,支持迭代、len()、bool()与in判断,迭代时按创建时间排序(src/textual/worker_manager.py)。
Worker 与创建它的 DOM 节点绑定(worker.node,见 src/textual/worker.py),这带来两重自动清理:
- 移除 Widget 或 pop 掉 Screen 时,其绑定的 Worker 任务自动被取消;
- 退出应用时,所有运行中的任务被取消。
这是发布说明强调的“有序关闭”机制的直接体现:Textual 原本就以“子节点先于父节点”的顺序关闭 Widget 任务,Worker API 搭上了这趟便车,因此你通常无需手工编写关闭逻辑。WorkerManager也提供了cancel_all()、cancel_group(node, group)、cancel_node(node)和wait_for_complete()等管理方法(src/textual/worker_manager.py),用于精细化控制。
Worker 状态机
每个 Worker 都有一个state属性,值为 WorkerState 枚举,标识其当前所处阶段:
| 值 | 含义 |
|---|---|
PENDING | Worker 已创建,尚未启动 |
RUNNING | 正在运行 |
CANCELLED | 已被取消,不再运行 |
ERROR | 抛出了异常,worker.error中保存异常对象 |
SUCCESS | 成功完成,worker.result中保存返回值 |
状态流转是单向的:PENDING→RUNNING→ 终态之一(CANCELLED/ERROR/SUCCESS),流程见仓库中的状态图 docs/images/workers/lifetime.excalidraw.svg。从源码看,state的 setter 在状态变化时会向创建节点投递StateChanged消息(src/textual/worker.py),这是下文事件机制的基石。此外Worker还提供is_running、is_finished、progress(配合update()/advance()报告完成度,见 src/textual/worker.py)等便捷属性。
Worker 事件:状态变化的通知机制
Worker 状态每次变化,都会向创建它的 Widget发送Worker.StateChanged事件(src/textual/worker.py)。通过定义on_worker_state_changed处理器即可监听,例如weather04.py中把生命周期事件全部记入 Textual 日志:
import httpx from rich.text import Text from textual import work from textual.app import App, ComposeResult from textual.containers import VerticalScroll from textual.widgets import Input, Static from textual.worker import Worker class WeatherApp(App): """App to display the current weather.""" CSS_PATH = "weather.tcss" def compose(self) -> ComposeResult: yield Input(placeholder="Enter a City") with VerticalScroll(id="weather-container"): yield Static(id="weather") async def on_input_changed(self, message: Input.Changed) -> None: """Called when the input changes""" self.update_weather(message.value) @work(exclusive=True) async def update_weather(self, city: str) -> None: """Update the weather for the given city.""" weather_widget = self.query_one("#weather", Static) if city: # Query the network API url = f"https://wttr.in/{city}" async with httpx.AsyncClient() as client: response = await client.get(url) weather = Text.from_ansi(response.text) weather_widget.update(weather) else: # No city, so just blank out the weather weather_widget.update("") def on_worker_state_changed(self, event: Worker.StateChanged) -> None: """Called when the worker state changes.""" self.log(event) if __name__ == "__main__": app = WeatherApp() app.run()用开发模式运行即可在 Textual 控制台中观察 PENDING → RUNNING → SUCCESS 的完整生命周期日志:
textual run weather04.py --dev由于StateChanged事件携带worker与state字段,这是在消息循环内安全读取worker.result/worker.error的标准姿势:既不需要await wait()阻塞界面,也避开了线程安全问题。
线程 Worker:处理不支持异步的库
前面所有示例都假设你使用 httpx 这类原生异步库;当第三方库不支持 async 时,可以设置thread=True把函数放进操作系统线程运行(run_worker内部通过loop.run_in_executor调度,见 src/textual/worker.py)。线程 Worker 的 API 与异步 Worker 完全一致,但有两个必须牢记的差异:
- 不要在线程 Worker 中直接操作 UI 或设置响应式变量——Textual 大部分函数都不是线程安全的。必须通过
App.call_from_thread把函数调度回主线程执行(src/textual/app.py,它内部用run_coroutine_threadsafe把回调投递到主事件循环);Widget.post_message是少数线程安全的方法,如果 Worker 需要多次更新界面,更好的做法是投递自定义消息,由消息处理器统一更新 UI 状态。 - 线程无法像协程那样被真正中断——你可以通过
get_current_worker()拿到当前 Worker 对象(src/textual/worker.py,基于 ContextVar 实现),然后手工检查worker.is_cancelled决定是否提前返回。
weather05.py用标准库urllib.request替换 httpx 完整演示了线程 Worker 的写法:
from urllib.parse import quote from urllib.request import Request, urlopen from rich.text import Text from textual import work from textual.app import App, ComposeResult from textual.containers import VerticalScroll from textual.widgets import Input, Static from textual.worker import Worker, get_current_worker class WeatherApp(App): """App to display the current weather.""" CSS_PATH = "weather.tcss" def compose(self) -> ComposeResult: yield Input(placeholder="Enter a City") with VerticalScroll(id="weather-container"): yield Static(id="weather") async def on_input_changed(self, message: Input.Changed) -> None: """Called when the input changes""" self.update_weather(message.value) @work(exclusive=True, thread=True) def update_weather(self, city: str) -> None: """Update the weather for the given city.""" weather_widget = self.query_one("#weather", Static) worker = get_current_worker() if city: # Query the network API url = f"https://wttr.in/{quote(city)}" request = Request(url) request.add_header("User-agent", "CURL") response_text = urlopen(request).read().decode("utf-8") weather = Text.from_ansi(response_text) if not worker.is_cancelled: self.call_from_thread(weather_widget.update, weather) else: # No city, so just blank out the weather if not worker.is_cancelled: self.call_from_thread(weather_widget.update, "") def on_worker_state_changed(self, event: Worker.StateChanged) -> None: """Called when the worker state changes.""" self.log(event) if __name__ == "__main__": app = WeatherApp() app.run()注意其中的安全模式:网络请求完成后先检查worker.is_cancelled,再通过self.call_from_thread(weather_widget.update, weather)把 UI 更新调度回主线程——这样既不会出现“过期响应覆盖新结果”,也避免了跨线程直接操作 Widget。
配套样式
所有天气示例共用同一个 weather.tcss:Input停靠在顶部占满宽度,#weather-container用height: 1fr撑满剩余空间并居中,天气内容区域自动尺寸并允许滚动,保证不同长度的 ANSI 天气文本都能正确展示。
小结
Textual 0.18.0 的 Worker API 把两件原本分散的事情统一到了一起:asyncio 任务与线程。run_worker与@work提供了一致的启动入口,Worker对象统一暴露状态、结果与取消能力,WorkerManager接管生命周期并把清理逻辑与既有的 DOM 树有序关闭机制对齐——这正是发布说明所称“解决 90% 的 Textual 应用并发问题”的底气所在。完整的进阶讲解(含事件驱动、devtools 调试等)可继续阅读 Worker 指南,对应源码则在 src/textual/worker.py、src/textual/worker_manager.py 与 src/textual/_work_decorator.py 中,仓库 tests/workers 目录下的测试用例(如 test_work_decorator.py、test_worker.py)覆盖了装饰器校验、状态流转、取消与异常路径,可作为深入研读的参考。
【免费下载链接】textualThe lean application framework for Python. Build sophisticated user interfaces with a simple Python API. Run your apps in the terminal and a web browser.项目地址: https://gitcode.com/gh_mirrors/te/textual
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考