ArkTS 并发编程实战:TaskPool 与 Worker 里如何安全地更新 UI
ArkTS 并发编程实战:TaskPool 与 Worker 里如何安全地更新 UI
前言
「子线程不能直接操作 UI」——这句话在 HarmonyOS 的技术社区里出现的频率,大概仅次于「@State 变化会触发刷新」。它是一句正确的结论,但在真实项目里,它常常被当成一个「知道了就能写对」的知识点。事实恰好相反:这句结论只解决了「要不要回主线程」的判断,完全没有回答「怎么回、回什么、回来之后怎么收尾」这三件真正会让代码写崩的事。
我在华为开发者论坛一篇题为《TaskPool / Worker 里想更新 UI 就抓瞎》的帖子下面,看到了这个问题的完整切面。提问者的诉求非常朴素:TaskPool 里不能直接改 @State,难道只能 postMessage 回来?Worker 太重,有没有更轻的方案?页面销毁了并发任务怎么取消?这三个问题其实分别对应并发编程里的通信、选型和生命周期,本来就不是一句话能答完的。
更有意思的是回复区。方向都对——「回主线程再赋值」「用 TaskPool 自带的进度通道」「aboutToDisappear 里取消」——但落地到代码,几乎每一份示例都有签名级别的问题:把 execute() 的返回值当成 Task 去 cancel()、把 sendData 当成注入进来的回调参数、把 isCanceled 写成 taskpool.isCanceled()、以为调了 cancel() 循环就会停下来。这些东西编译器会替你挡掉一部分,剩下的一部分会在真机上以「点了取消进度条还在爬」「页面退了还在写状态」「连点两次结果串台」的形式暴露出来,而且很难定位到根因。
还有一层认知债被反复复制:很多人从 React / Flutter 迁过来,下意识觉得「setState 谁都能调,框架内部保证线程安全」。在 ArkTS 里这个前提根本不存在,ArkTS 的状态变量赋值就是一次直接的字段写入,没有任何队列、调度器或锁在中间兜底。把这个前提纠正过来,后面所有的写法选择都会变得顺理成章。
这篇文字做四件事:先把「不能操作 UI」这条铁律的边界划清楚,说清它到底约束了什么、没约束什么;再逐条纠正那些会让方案跑不起来的 API 误用;然后把 cancel() 与 sendData 的真实语义摊开——它们的行为和大多数人的直觉相反,这正是误判「框架有 bug」的主要来源;最后给出两份可以直接抄进项目的骨架:一份是 TaskPool 的进度上报与协作式取消,一份是 Worker 长驻流式任务的完整收发。
问题描述
原始诉求可以浓缩成一句话:在并发任务中做一件需要持续十几秒的同步计算,期间进度条要跟着走,用户中途退出页面时任务要停,回来后结果不能写进已经销毁的组件。
拆开看,这里有四个彼此独立的技术难点,每一个单独拿出来都不算复杂,但凑在一起就会出现组合爆炸:
- 结果怎么回来。 子线程算完的
number,要通过什么通道回到 UI 线程?await taskpool.execute()之后直接赋值,是不是就够了? - 中间进度怎么回来。 最终结果只有一个,进度可能有几百次。每一次都走
execute的返回值显然不可能,那走什么? - 中途不想算了怎么办。 用户点了取消,或者页面被销毁了,任务要怎么停?
- 停了之后,那些已经在路上的数据怎么办。 取消是一个动作,而数据回传是持续的流;两者之间存在时间差,这个时间差里发生的事才是 bug 的温床。
社区里围绕这四个问题给出的方案,大致有三类,每一类都解决了一部分,也都在别的地方留了坑。
第一类是「Emitter 万能论」:主线程用 emitter.on 订阅,并发函数里 emitter.emit 触发,页面销毁时 off 掉。这条路径本身是官方认可的能力,问题出在工程细节——订阅写在 onClick 里,连点 N 次就有 N 个订阅;off(eventId) 取消的是该事件 ID 下的所有订阅,于是两个页面用了同一个裸 ID(帖子里的例子就是 "1")之后,B 页面注销会把 A 页面的订阅一起拔掉,这种互相踩的现象极难排查。
第二类是「sendData 当成回调注入」:写成 taskpool.execute(calcWithProgress, (data) => { this.progress = data.progress }),认为第三个参数会被框架识别成回调。这个写法不成立——execute() 的参数列表里,除第一个函数引用之外的所有东西都是任务入参,会被序列化后传给并发函数。它不会被识别成任何回调,运行时只会得到一个类型不匹配或序列化失败。
第三类是「保存 execute() 的返回值来取消」:this.task = taskpool.execute(f, 1),然后在 aboutToDisappear 里 taskpool.cancel(this.task)。这段代码过不了编译,因为 execute() 返回的是 Promise,而 cancel() 的形参类型是 taskpool.Task。要能取消,必须显式 new taskpool.Task(...),把 Task 对象留下来。
把这三类方案里的错误集中起来,就得到了本文要逐条纠正的清单:
taskpool.execute()的返回值是Promise,不可传给cancel();sendData/isCanceled是taskpool.Task的静态方法,既不是taskpool上的方法,也不是注入进任务函数的回调;sendData只能在并发函数自己的执行线程内被调用,放进setTimeout、Promise.then、事件回调里调用都会失败;onReceiveData()必须在taskpool.execute()之前注册,顺序反了第一条消息必然丢;cancel()是协作式取消:取消后任务的执行体不会被打断,进度仍可能继续回传;sendData不保证顺序、不保证不丢,不适合传最终结果;- emitter 存在重复订阅与裸 eventId 互踩的隐患,需要统一的 channel 管理。
这些问题里,前四条是「写错了会直接失败」,后三条是「写对了也可能不符合预期」。后者更危险,因为它不会报错,只会让你在真机上怀疑人生。
细节解析
一、先把「子线程不能操作 UI」这条铁律的边界划清楚
这条规则经常被两种极端理解同时污染:一种认为「所有跨线程的 UI 交互都必须自己手写消息队列」,另一种认为「只要最终赋值那一行在主线程执行就行,中间无所谓」。两种都不准确。
准确的表述是:ArkTS 的 UI 实例、状态变量以及驱动它们刷新的渲染管线,全部与 UI 线程绑定;任何线程只要触碰这条链路,都必须在 UI 线程上执行,且必须通过框架提供的回传通道。
拆开说,被约束的是三样东西。
第一样是状态变量的写入。@State、@Prop、@Link、@Observed 这些装饰器背后是一套「读时收集依赖、写时打脏标记」的机制。打脏标记要往 UI 实例持有的依赖表里写数据,而那张表既没有加锁,也没有任何内存屏障语义,它从设计上就假设只有一个线程访问它。从子线程写进去,轻则脏标记丢失(状态变了但界面不刷新),重则破坏依赖表结构导致渲染异常甚至崩溃。
第二样是渲染相关的上下文。UIContext、帧回调(vsync)、animateTo 这类能力的注册,都要求当前线程存在有效的 UI 上下文。子线程里没有这个上下文,调用会直接被框架拒绝。这一点解释了一个常见现象:为什么 animateTo 在并发函数里连报错都显得莫名其妙。
第三样是编译期的可见性。被 @Concurrent 标记的函数会被编译器做作用域裁剪:它不能访问组件实例(拿不到 this)、不能访问 UI 相关对象、不能引用模块级的可变变量,入参和返回值还必须可序列化。所以「在并发函数里写 this.progress = 50」这件事根本走不到运行时——编译器在更早的地方就把它拦住了。
那没有被约束的是什么?是「计算本身」。子线程可以做大数组遍历、加解密、JSON 解析、图像像素处理、数据库查询、网络请求;可以产出中间结果、可以调用 sendData、可以 emit 事件。换句话说,约束在「写 UI」这一端,不在「算」这一端。把这条边界记牢,就不会再出现「为了更新一个进度条,把整个计算过程都搬到主线程」这种本末倒置的写法。
顺便澄清一个高频误解:很多从 React 或 Flutter 迁过来的开发者会问「ArkTS 有没有线程安全的 setState」。没有,而且不是「暂时没有」。React 的 setState 能跨线程调,前提是它有一个明确的批处理调度器——调用只是把更新入队,真正的合并与渲染由调度器在渲染线程完成。ArkTS 的状态赋值语义是「直接写字段 + 打脏标记」,中间没有队列、没有调度层、没有批次合并。缺失的不是一个 API,而是整个调度层。所以在 ArkTS 里追求「setState 式」的线程安全写法,方向就是错的。
还有一个值得知道的边界松动:HarmonyOS 后续版本引入了 @Sendable 装饰的可跨线程共享对象,配合主线程侧的观察者机制(@ObservedV2 / @Trace 与 makeObserved),子线程修改共享对象的属性、主线程观测到变化并驱动刷新,是可以做到的。但请注意这条路径的实质:驱动刷新的仍然是主线程侧的观察者,不是子线程自己去碰渲染管线。它把「数据的跨线程共享」和「刷新的触发」这两件事分开处理了,铁律本身并没有被推翻。
二、三个会让方案直接跑不起来的 API 误用
误用一:把 execute() 的返回值当成 Task。
taskpool.execute() 有两组重载:execute(func, ...args) 和 execute(task, priority?)。前者的返回值是 Promise<Object>——计算结果的异步句柄;后者的返回值同样是 Promise<Object>,而不是 Task。taskpool.cancel() 接受的是 taskpool.Task 对象,两者类型完全不同,编译器会直接报错。
正确的做法是把「创建任务」和「提交任务」分开:
const task: taskpool.Task = new taskpool.Task(scanWithProgress, TOTAL);
this.activeTask = task; // 先留住 Task 对象
const result: number = await taskpool.execute(task) as number;
这个改动看着只是多写了一行,但它是「任务可取消」的前提。先 execute(func, ...args) 拿到 Promise 再想办法取消,是没有路径的——Task 对象根本没被创建出来,cancel() 无从下手。
误用二:sendData / isCanceled 的宿主对象搞错了。
这两个是 taskpool.Task 的静态方法,调用形式是 taskpool.Task.sendData(...) 和 taskpool.Task.isCanceled()。写成 taskpool.sendData(...) / taskpool.isCanceled() 会直接编译失败。这一点在复制粘贴的帖子里传播得非常广。
误用三:把 sendData 当成注入进来的回调参数。
这是最有迷惑性的一类。有人写成:
@Concurrent
function calcWithProgress(sendData: taskpool.SendDataCallback): void {
for (let i = 1; i <= 100; i++) {
sendData({ progress: i });
}
}
taskpool.execute(calcWithProgress, (data) => { this.progress = data.progress; });
看起来「能跑」,实质上完全不成立。execute() 在函数引用之后的所有参数都是任务入参,会被序列化后一并传给并发函数。传一个箭头函数进去,序列化阶段就会出问题;即使侥幸通过,那个参数也只是个普通入参,不会携带任何回传能力。正确的对应关系是:并发函数内用 taskpool.Task.sendData() 发送,宿主侧用 task.onReceiveData() 接收,两者通过同一个 Task 对象绑定。
还有两条硬纪律:sendData 只能在并发函数自己的执行线程内调用,把它挪进 setTimeout、Promise.then、事件回调这类「脱离当前执行流」的位置就会失败;onReceiveData() 必须先注册、再提交,也就是先 task.onReceiveData(cb) 再 taskpool.execute(task),顺序反了的话,第一条消息必然在还没有接收者的时候被丢掉。
三、cancel() 是协作式取消:它改的是调用方视角,不是执行体
这是整篇文章里最容易被误判成「框架 bug」的一条。
taskpool.cancel(task) 做的事情是:把任务标记为已取消,让 execute() 返回的 Promise 以拒绝结束,并把最终结果丢弃。它不会打断并发函数里正在执行的循环,也不会去杀线程。任务函数会继续跑到自己的自然退出点,期间 sendData 照样继续发。
所以「点了取消,进度条还一路爬到 100%」是符合设计的行为,不是异常。想让计算真正停下来,唯一的办法是在并发函数内部按分段边界主动检查:
@Concurrent
function scanWithProgress(total: number): number {
for (let i: number = 0; i < total; i++) {
if (i % 65536 === 0) {
if (taskpool.Task.isCanceled()) {
return -1; // 协作式退出点
}
taskpool.Task.sendData(Math.floor(i * 100 / total));
}
}
return total;
}
这里有两个工程上的取舍值得说明。检查点的密度决定了取消的响应速度:检查点太稀疏,用户点了取消要等好几秒才停;每个循环都查一次,又会把 isCanceled() 的调用开销摊到每一次迭代上。经验做法是按数据分段来,比如每处理 64K 条记录检查一次,既能把响应延迟压到毫秒级,又不会影响主循环性能。
取消的可靠性还跟任务状态有关。 还在排队、尚未开始执行的任务,可以直接取消掉;已经在运行中的任务,取消能否生效与目标 API 版本有关,较低版本对运行中任务的取消可能失败并抛出 BusinessError(10200016 就属于这类「任务正在执行,无法取消」的情形)。所以 cancel() 一定要包 try / catch,而且永远不要把「取消成功」当成「不会再有回调」的保证。
这就是为什么每个页面上都需要一个看起来有点多余的东西——代次(generation)校验。页面侧维护一个自增的 gen,每发起一次新任务就 ++gen 并记下本次的 myGen;所有异步回调(进度、结果、异常)进入时先比一次 myGen === this.gen,不相等就直接丢弃。再配合一个 alive 存活标记(aboutToAppear 置 true、aboutToDisappear 置 false),就同时挡住了三种脏数据:
- 页面已销毁后到达的进度与结果(
alive为false); - 上一次任务在取消后继续发出来的进度(
gen已递增); - 快速连点「开始」时,前一次任务迟到的结果覆盖后一次的结果(
gen已递增)。
需要强调的是,alive 和 gen 校验不能省,哪怕已经正确调用了 cancel()。它们和「取消」是两条独立的防线:取消负责「尽量让计算停下来」,代次校验负责「无论计算停没停下来,脏数据都不许进 UI」。
还有一个容易踩的细节:页面隐藏不等于页面销毁。如果任务很长、且用户可能切到其他页面再回来,用 aboutToDisappear 就取消任务,会导致用户回来看到进度归零、计算白跑。更稳妥的策略是在 onPageHide 里只做「停止回传进度」,在 aboutToDisappear 里才真正取消;或者在退到后台时暂停 UI 侧接收、把 Task 保留在成员变量里,回到前台再恢复接收。这属于产品策略层面的事,但必须在写代码前想清楚,因为它决定了你把取消逻辑挂到哪个生命周期上。
四、sendData 的语义边界:它是进度通道,不是结果通道
sendData 走了底层线程池的消息队列,它有几个必须知道的特性:不保证顺序、不保证不丢、不保证一定送达。它适合承载「进度百分比」这类可容忍丢失、且丢了也无所谓的中间态数据。它不适合承载最终结果、不适合承载必须严格有序的序列、更不适合当作可靠队列使用。
最终结果永远走 execute() 的 resolve 值。这是唯一能保证「任务正常跑完 → 结果一定正确送达」的通道,因为它的语义是「任务完成」,而不是「发了一条消息」。
还有一个实测经验值得记下来:回传频率要节制。每 1% 发一次进度,一个跑十几秒的任务可能每秒写几十次 @State;进度条确实顺滑,但页面的点击响应会明显变钝,因为主线程被高频的重绘请求占住了。把策略改成「累计变化满 3%~5% 才写一次状态变量」,观感上没有任何差别,重绘压力却能降一个量级。如果对流畅度有更高要求,进一步的做法是让进度只在文本上高频更新、而进度条本身走离散的分段,或者干脆用动画插值把两次真实进度之间的过渡补齐,把「真实的进度点」和「视觉上的连续性」解耦。
五、用 emitter 做回传通道的两个后遗症与统一 channel 管理
emitter 是官方支持的线程间通信能力,emit 可以在子线程调用,on 在 UI 线程注册,跨线程传递的数据同样要可序列化。把它当回传通道本身没问题,问题在后遗症。
后遗症一:重复订阅。 emitter.on('1', cb) 如果写在 onClick 里,连点 N 次就产生 N 个订阅,每个回调都持有组件引用。页面销毁后这些引用仍可能被调用(往已销毁组件写状态),同时内存持续增长。注销有两种写法,语义不同,必须分清:
emitter.off(eventId, handler):精确注销指定 handler,前提是注销时传的那个 handler 引用,和订阅时是同一个。所以匿名箭头函数写在on里就很难精确注销,必须把 handler 提到成员变量里保存引用。emitter.off(eventId):注销该 eventId 下的所有订阅。方便,但杀伤面大。
后遗症二:裸 eventId 互踩。 这正好是上一条的反面:如果两个页面都订阅了 "1" 这个裸 ID,A 页面走 off("1") 就会把 B 页面的订阅一起拔掉,B 页面从此收不到任何消息,而且不报错。
结论是:eventId 必须有归属。推荐用「前缀 + 作用域 + 主题」的复合 ID,例如 harmony.demo.channel.<pageScope>.<topic>,再用一个薄封装把「订阅 / 注销 / 同主题去重」收口到一处,避免 on 和 off 散落在各处。封装要守住三条:同一 scope + topic 重复订阅时先释放旧订阅;注销时用保存下来的 handler 引用做精确注销;页面销毁时按 scope 批量释放。
那 emitter 和 sendData 该怎么分工?我的建议很明确:同一个任务内部的进度与结果,优先用 sendData / onReceiveData 或 execute() 的返回值,因为它们天然绑定任务归属,注销和去重都不需要你操心(Task 对象被回收,通道就跟着消失)。emitter 更适合「跨模块、发布订阅式」的通知——全局登录态变化、下载完成广播、跨页面的数据刷新这类「一个事件有多个、且不固定的消费者」的场景。
六、TaskPool 还是 Worker:选型判据
TaskPool 是在 Worker 之上实现的调度器和线程池:你把任务丢进去,线程的创建、复用、回收由框架管理,你不需要关心生命周期。Worker 则是一个你自己显式创建、需要自己 terminate 的长驻线程,通信只能靠消息。两者的能力不是「谁更强」,而是「谁更贴合场景」。
先说什么不值得进池。 @Concurrent 函数的入参和返回值都要走序列化 / 反序列化,调度本身也有成本。经验阈值是:单次同步计算持续 1ms 以上才值得丢进 TaskPool,低于这个量级,序列化开销很可能比省下来的时间还多。另外,IO 操作(网络请求、文件读写、数据库查询)本身就是异步 API,await 一下就行,不要为了「异步」再包一层 TaskPool——那只是把一个异步包装成另一个异步,还额外付了序列化的钱。
需要吞吐时按核数分片。 真正吃 CPU 的批量计算,拆成若干个分片并发提交,用 Promise.all 收敛:
const shards: number = 4;
const step: number = Math.ceil(TOTAL / shards);
const jobs: Promise<number>[] = [];
for (let i: number = 0; i < shards; i++) {
const task: taskpool.Task =
new taskpool.Task(rangeSum, i * step, Math.min((i + 1) * step, TOTAL));
jobs.push(taskpool.execute(task) as Promise<number>);
}
const parts: number[] = await Promise.all(jobs);
const total: number = parts.reduce((a: number, b: number) => a + b, 0);
注意别一次提交几百个任务——线程池扩容有上限,多余任务只会在队列里排队,收益递减。
什么场景必须用 Worker。 判据只有一条:线程是否需要长期驻留并维护自己的状态。典型的比如:长连接的心跳与消息解包、需要在子线程内维护一个持续增长的数据结构的轮询任务、需要持续消费某个数据流的处理管线。这类任务的共同点是「不是一个有明确起止的计算」,把它硬塞进 TaskPool 会变成「在并发函数里开一个死循环」,既拿不到结果返回值,也没有自然的完成时机,反而更难管理。
Worker 的两个坑必须提前知道。 第一,terminate() 是异步的:调用之后线程不会立刻消失,仍可能再投递一两条消息过来。所以宿主侧的 on('message') 回调里同样要做存活与代次校验,不能因为「已经 terminate 了」就假设不会再被调用。第二,Worker 的 postMessage 默认走结构化克隆(深拷贝),大对象来回传的成本不低;如果消息体很大,应该考虑用 Transferable(ArrayBuffer 的所有权转移,零拷贝)或 @Sendable 共享对象来降低开销。子线程里要用的对象、跨线程传递的对象,也不要指望它们是同一个引用。
示例代码
下面给三份代码。第一份是 TaskPool 的完整骨架,覆盖进度上报、协作式取消、代次校验、销毁清理;第二份是统一的 emitter channel 封装,解决重复订阅与 ID 互踩;第三份是 Worker 长驻流式任务,包含宿主页面和线程文件两个部分。
示例一:TaskPool 进度上报 + 协作式取消 + 代次校验
并发函数与页面放在同一个 .ets 文件里(也可拆到单独的 util 文件,只要用 @Concurrent 且是顶层函数)。
// entry/ets/pages/TaskPoolProgressPage.ets
import { taskpool } from '@kit.ArkTS';
import { BusinessError } from '@kit.BasicServicesKit';
const TOTAL: number = 20_000_000;
const CHECK_INTERVAL: number = 65536; // 每处理 64K 条检查一次取消
const REPORT_STEP: number = 3; // 进度累计变化满 3% 才回传一次
// 顶层函数,@Concurrent 标记,不能访问组件实例、UI、模块级可变变量
@Concurrent
function scanWithProgress(total: number): number {
let hits: number = 0;
let lastReported: number = -1;
for (let i: number = 0; i < total; i++) {
if (i % CHECK_INTERVAL === 0) {
// 静态方法调用,用于协作式退出
if (taskpool.Task.isCanceled()) {
return -1;
}
const percent: number = Math.floor(i * 100 / total);
if (percent - lastReported >= REPORT_STEP) {
lastReported = percent;
// 只能在本并发函数自己的执行线程内调用
taskpool.Task.sendData(percent);
}
}
if (i % 13 === 0) {
hits++;
}
}
return hits;
}
@Entry
@Component
struct TaskPoolProgressPage {
@State percent: number = 0;
@State status: string = '待开始';
@State running: boolean = false;
private activeTask?: taskpool.Task;
private generation: number = 0;
private alive: boolean = false;
aboutToAppear(): void {
this.alive = true;
}
aboutToDisappear(): void {
this.alive = false;
this.generation++; // 让所有在途回调失效
this.cancelActiveTask();
}
private cancelActiveTask(): void {
const task: taskpool.Task | undefined = this.activeTask;
this.activeTask = undefined;
if (task === undefined) {
return;
}
try {
taskpool.cancel(task);
} catch (error) {
const err: BusinessError = error as BusinessError;
if (err.code !== 10200016) {
console.warn(`cancel task failed: ${err.code} ${err.message}`);
}
}
}
async start(): Promise<void> {
if (this.running) {
return;
}
this.cancelActiveTask();
const gen: number = ++this.generation;
const task: taskpool.Task = new taskpool.Task(scanWithProgress, TOTAL);
// 必须在 execute 之前注册,否则第一条进度会丢
task.onReceiveData((data: Object) => {
if (!this.alive || gen !== this.generation) {
return; // 销毁后或过期任务的进度直接丢弃
}
this.percent = data as number;
this.status = `扫描中 ${this.percent}%`;
});
this.activeTask = task;
this.running = true;
this.percent = 0;
this.status = '扫描中 0%';
try {
// 最终结果只认 execute 的 resolve 值,不走 sendData
const hits: number = await taskpool.execute(task) as number;
if (!this.alive || gen !== this.generation) {
return;
}
if (hits < 0) {
this.status = '已取消';
} else {
this.percent = 100;
this.status = `完成,命中 ${hits} 条`;
}
} catch (error) {
const err: BusinessError = error as BusinessError;
if (this.alive && gen === this.generation) {
this.status = `任务结束:${err.code}`;
}
} finally {
if (this.activeTask === task) {
this.activeTask = undefined;
}
if (this.alive && gen === this.generation) {
this.running = false;
}
}
}
build() {
Column({ space: 16 }) {
Text(this.status)
.fontSize(16)
.width('100%')
Progress({ value: this.percent, total: 100 })
.width('100%')
Row({ space: 12 }) {
Button('开始')
.enabled(!this.running)
.onClick(() => {
this.start();
})
Button('取消')
.enabled(this.running)
.onClick(() => {
// 先让在途数据失效,再提交取消
this.generation++;
this.cancelActiveTask();
this.running = false;
this.status = '已取消(任务会在下一个检查点退出)';
})
}
.width('100%')
}
.width('100%')
.padding(24)
}
}
这份代码里有几个地方是刻意写出来的,不能省:cancelActiveTask() 里的 try / catch(吞掉「任务正在执行无法取消」);onReceiveData 里对 alive 和 gen 的双重校验(取消后仍可能收到进度);finally 里按 Task 引用相等来判断是否清空(防止误清掉新任务);以及取消按钮里先 ++generation 再 cancel 的顺序。
示例二:统一的 emitter channel 管理
这个封装解决三件事:复合 eventId 避免互踩、同主题重复订阅先去重、按作用域批量释放。
// entry/ets/common/EventChannel.ets
import { emitter } from '@kit.BasicServicesKit';
type ChannelHandler = (data: emitter.EventData) => void;
export class EventChannel {
private static readonly PREFIX: string = 'harmony.demo.channel';
private handlers: Map<string, ChannelHandler> = new Map<string, ChannelHandler>();
private keyOf(scope: string, topic: string): string {
return `${EventChannel.PREFIX}.${scope}.${topic}`;
}
// 订阅:同一 scope + topic 重复订阅时,先释放旧订阅,避免回调叠加
subscribe(scope: string, topic: string, handler: ChannelHandler): void {
const id: string = this.keyOf(scope, topic);
this.release(scope, topic);
emitter.on(id, handler);
this.handlers.set(id, handler);
}
// 发布:可在子线程调用,payload 必须可序列化
publish(scope: string, topic: string, payload: Object): void {
const eventData: emitter.EventData = { data: payload };
emitter.emit(this.keyOf(scope, topic), eventData);
}
// 精确注销:使用与订阅时相同的 handler 引用
release(scope: string, topic: string): void {
const id: string = this.keyOf(scope, topic);
const handler: ChannelHandler | undefined = this.handlers.get(id);
if (handler !== undefined) {
emitter.off(id, handler);
this.handlers.delete(id);
}
}
// 按作用域批量释放,供 aboutToDisappear 调用
releaseScope(scope: string): void {
const prefix: string = `${EventChannel.PREFIX}.${scope}.`;
this.handlers.forEach((handler: ChannelHandler, id: string) => {
if (id.startsWith(prefix)) {
emitter.off(id, handler);
this.handlers.delete(id);
}
});
}
}
export const eventChannel: EventChannel = new EventChannel();
页面侧的使用方式很直接,订阅与释放成对出现,scope 用页面自己的标识:
import { eventChannel } from '../common/EventChannel';
import { emitter } from '@kit.BasicServicesKit';
@Entry
@Component
struct NotifyPage {
@State message: string = '等待广播';
private readonly scope: string = 'NotifyPage';
aboutToAppear(): void {
eventChannel.subscribe(this.scope, 'download', (data: emitter.EventData) => {
const payload: Record<string, Object> = data.data as Record<string, Object>;
this.message = `下载完成:${payload['name'] as string}`;
});
}
aboutToDisappear(): void {
eventChannel.releaseScope(this.scope);
}
build() {
Column({ space: 12 }) {
Text(this.message)
Button('模拟完成广播')
.onClick(() => {
eventChannel.publish(this.scope, 'download', { name: 'report.pdf' });
})
}
.padding(24)
}
}
示例三:Worker 长驻流式任务
Worker 需要一个独立的线程文件。宿主页面负责创建实例、注册监听、发送指令,以及退出时的清理;线程文件负责持续消费数据并回传进度与结果。
// entry/ets/workers/StreamWorker.ets
import { worker } from '@kit.ArkTS';
interface HostCommand {
cmd: string;
total: number;
}
interface WorkerEvent {
type: string;
seq: number;
percent: number;
hits: number;
detail: string;
}
const workerPort: worker.ThreadWorkerGlobalScope = worker.workerPort;
let running: boolean = false;
let canceled: boolean = false;
function emit(type: string, seq: number, percent: number, hits: number, detail: string): void {
const event: WorkerEvent = {
type: type,
seq: seq,
percent: percent,
hits: hits,
detail: detail
};
workerPort.postMessage(event);
}
function run(total: number): void {
if (running) {
emit('rejected', 0, 0, 0, '已有任务在执行');
return;
}
running = true;
canceled = false;
const batch: number = Math.max(1, Math.floor(total / 100));
let hits: number = 0;
let seq: number = 0;
for (let i: number = 0; i < total; i++) {
// 协作式停止:宿主侧发 stop 指令后在这里退出,保留线程以便复用
if (canceled) {
running = false;
emit('canceled', seq, Math.floor(i * 100 / total), hits, '已停止');
return;
}
if (i % 29 === 0) {
hits++;
}
if (i % batch === 0) {
seq++;
emit('progress', seq, Math.floor(i * 100 / total), hits, '处理中');
}
}
running = false;
emit('done', seq + 1, 100, hits, '全部处理完成');
}
workerPort.on('message', (e: MessageEvents) => {
const command: HostCommand = e.data as HostCommand;
if (command.cmd === 'start') {
run(command.total);
} else if (command.cmd === 'stop') {
canceled = true;
}
});
宿主页面对应的实现如下。注意 terminate() 之后仍可能收到消息,所以回调里保留了存活校验。
// entry/ets/pages/WorkerStreamPage.ets
import { worker } from '@kit.ArkTS';
import { BusinessError } from '@kit.BasicServicesKit';
interface WorkerEvent {
type: string;
seq: number;
percent: number;
hits: number;
detail: string;
}
@Entry
@Component
struct WorkerStreamPage {
@State percent: number = 0;
@State status: string = '待开始';
@State running: boolean = false;
private workerInstance: worker.ThreadWorker | null = null;
private generation: number = 0;
private alive: boolean = false;
aboutToAppear(): void {
this.alive = true;
}
aboutToDisappear(): void {
this.alive = false;
this.generation++;
this.stopWorker();
}
private ensureWorker(): worker.ThreadWorker {
const current: worker.ThreadWorker | null = this.workerInstance;
if (current !== null) {
return current;
}
const instance: worker.ThreadWorker =
new worker.ThreadWorker('entry/ets/workers/StreamWorker.ets', { name: 'stream-worker' });
const gen: number = this.generation;
instance.on('message', (e: MessageEvents) => {
// terminate 是异步的,回调仍可能到达,必须校验
if (!this.alive || gen !== this.generation) {
return;
}
const event: WorkerEvent = e.data as WorkerEvent;
this.percent = event.percent;
if (event.type === 'progress') {
this.status = `处理中 ${event.percent}%,累计命中 ${event.hits}`;
} else if (event.type === 'done') {
this.status = `完成,累计命中 ${event.hits}`;
this.running = false;
} else if (event.type === 'canceled') {
this.status = '已停止';
this.running = false;
} else if (event.type === 'rejected') {
this.status = event.detail;
}
console.info(`worker seq=${event.seq} type=${event.type}`);
});
instance.on('error', (err: BusinessError) => {
if (!this.alive || gen !== this.generation) {
return;
}
this.status = `线程异常:${err.code}`;
this.running = false;
console.error(`worker error: ${err.code} ${err.message}`);
});
this.workerInstance = instance;
return instance;
}
private startWorkerTask(): void {
if (this.running) {
return;
}
this.generation++;
const instance: worker.ThreadWorker = this.ensureWorker();
this.running = true;
this.percent = 0;
this.status = '处理中 0%';
instance.postMessage({ cmd: 'start', total: 12_000_000 });
}
private requestStop(): void {
const instance: worker.ThreadWorker | null = this.workerInstance;
if (instance === null) {
return;
}
// 优先走协作式停止,线程可复用;不销毁线程
instance.postMessage({ cmd: 'stop', total: 0 });
this.running = false;
this.status = '停止指令已发送';
}
private stopWorker(): void {
const instance: worker.ThreadWorker | null = this.workerInstance;
this.workerInstance = null;
if (instance === null) {
return;
}
try {
instance.terminate();
} catch (error) {
const err: BusinessError = error as BusinessError;
console.warn(`terminate worker failed: ${err.code}`);
}
}
build() {
Column({ space: 16 }) {
Text(this.status)
.fontSize(16)
.width('100%')
Progress({ value: this.percent, total: 100 })
.width('100%')
Row({ space: 12 }) {
Button('启动')
.enabled(!this.running)
.onClick(() => {
this.startWorkerTask();
})
Button('停止')
.enabled(this.running)
.onClick(() => {
this.requestStop();
})
Button('销毁线程')
.onClick(() => {
this.stopWorker();
this.running = false;
this.status = '线程已销毁';
})
}
.width('100%')
}
.width('100%')
.padding(24)
}
}
这里把「停止」和「销毁线程」拆成了两个按钮,就是为了体现前面的结论:长驻线程上的任务取消,应该优先走协作式停止(发指令 + 线程内检查标志位),而不是一上来就 terminate。terminate 会让整个线程连同它维护的状态一起消失,重建成本高;只有在页面真正销毁、或者线程已经进入不可恢复状态时,才应该销毁它。
总结
把上面的内容压缩成一份可以贴在工位上的清单。
关于边界。 「子线程不能操作 UI」约束的是状态变量写入、渲染上下文的使用、以及编译期的可见性;它不约束计算本身。ArkTS 里没有、也不会有 setState 那种「调用入队、调度器兜底」的线程安全赋值,因为状态赋值的语义就是直接写字段,中间没有任何调度层。跨线程共享对象 + 主线程观察者的组合可以把边界往外推一点,但驱动刷新的仍然是主线程侧的观察者。
关于 API。 taskpool.execute() 返回的是 Promise,不是 Task;要能取消,必须先 new taskpool.Task(...) 并留住对象。sendData 和 isCanceled 都是 taskpool.Task 的静态方法,既不是 taskpool 上的方法,也不是注入进任务函数的回调。execute() 在函数引用之后的参数全部是任务入参,传回调进去不会生效。sendData 只能在并发函数自己的执行线程内调用,onReceiveData() 必须先于 execute() 注册。
关于取消。 cancel() 是协作式的:它能让 execute() 的 Promise 拒绝、结果被丢弃,但不会打断执行体,进度仍可能继续回传。真正停下来只能靠在并发函数内部按分段边界检查 taskpool.Task.isCanceled() 并主动 return。排队中的任务可以直接取消,运行中的任务取消能否生效与 API 版本有关,低版本可能抛 10200016,所以 try / catch 不能省。无论如何,cancel() 都不能替代代次校验。
关于通道。 最终结果走 execute() 的 resolve 值;中间进度走 sendData / onReceiveData,并做好节流(累计变化满 3%~5% 才写状态变量),因为它不保证顺序、不保证不丢。emitter 更适合跨模块的发布订阅式通知,用复合 eventId 保证归属,用同一个 handler 引用做精确注销,页面销毁时按 scope 批量释放。
关于选型。 单次同步计算超过 1ms 才值得进 TaskPool;纯 IO 直接 await;需要吞吐时按核数分片并用 Promise.all 收敛,不要一次提交几百个任务。线程需要长期驻留并维护自身状态时用 Worker,记得 terminate() 是异步的,回调里仍要校验存活与代次。
关于生命周期。 页面隐藏不等于销毁,取消逻辑挂在哪个生命周期上要在写代码前定下来。无论选哪条路径,alive 标记 + generation 代次这两道校验都要保留——它们是并发任务与 UI 之间最后一道、也是最可靠的一道防线。取消负责尽力而为,代次校验负责万无一失。
真正会出问题的从来不是「不知道要回主线程」,而是回主线程的路上丢的那几帧、多写的那几次状态、以及在页面已经销毁之后才姗姗来迟的那一条进度。把这几个环节管住,TaskPool 和 Worker 用起来其实比想象中干净得多。
更多推荐

所有评论(0)