ArkTS 并发编程实战:TaskPool 与 Worker 里如何安全地更新 UI

前言

「子线程不能直接操作 UI」——这句话在 HarmonyOS 的技术社区里出现的频率,大概仅次于「@State 变化会触发刷新」。它是一句正确的结论,但在真实项目里,它常常被当成一个「知道了就能写对」的知识点。事实恰好相反:这句结论只解决了「要不要回主线程」的判断,完全没有回答「怎么回回什么回来之后怎么收尾」这三件真正会让代码写崩的事。

我在华为开发者论坛一篇题为《TaskPool / Worker 里想更新 UI 就抓瞎》的帖子下面,看到了这个问题的完整切面。提问者的诉求非常朴素:TaskPool 里不能直接改 @State,难道只能 postMessage 回来?Worker 太重,有没有更轻的方案?页面销毁了并发任务怎么取消?这三个问题其实分别对应并发编程里的通信选型生命周期,本来就不是一句话能答完的。

更有意思的是回复区。方向都对——「回主线程再赋值」「用 TaskPool 自带的进度通道」「aboutToDisappear 里取消」——但落地到代码,几乎每一份示例都有签名级别的问题:把 execute() 的返回值当成 Taskcancel()、把 sendData 当成注入进来的回调参数、把 isCanceled 写成 taskpool.isCanceled()、以为调了 cancel() 循环就会停下来。这些东西编译器会替你挡掉一部分,剩下的一部分会在真机上以「点了取消进度条还在爬」「页面退了还在写状态」「连点两次结果串台」的形式暴露出来,而且很难定位到根因。

还有一层认知债被反复复制:很多人从 React / Flutter 迁过来,下意识觉得「setState 谁都能调,框架内部保证线程安全」。在 ArkTS 里这个前提根本不存在,ArkTS 的状态变量赋值就是一次直接的字段写入,没有任何队列、调度器或锁在中间兜底。把这个前提纠正过来,后面所有的写法选择都会变得顺理成章。

这篇文字做四件事:先把「不能操作 UI」这条铁律的边界划清楚,说清它到底约束了什么、没约束什么;再逐条纠正那些会让方案跑不起来的 API 误用;然后把 cancel()sendData 的真实语义摊开——它们的行为和大多数人的直觉相反,这正是误判「框架有 bug」的主要来源;最后给出两份可以直接抄进项目的骨架:一份是 TaskPool 的进度上报与协作式取消,一份是 Worker 长驻流式任务的完整收发。

问题描述

原始诉求可以浓缩成一句话:在并发任务中做一件需要持续十几秒的同步计算,期间进度条要跟着走,用户中途退出页面时任务要停,回来后结果不能写进已经销毁的组件。

拆开看,这里有四个彼此独立的技术难点,每一个单独拿出来都不算复杂,但凑在一起就会出现组合爆炸:

  1. 结果怎么回来。 子线程算完的 number,要通过什么通道回到 UI 线程?await taskpool.execute() 之后直接赋值,是不是就够了?
  2. 中间进度怎么回来。 最终结果只有一个,进度可能有几百次。每一次都走 execute 的返回值显然不可能,那走什么?
  3. 中途不想算了怎么办。 用户点了取消,或者页面被销毁了,任务要怎么停?
  4. 停了之后,那些已经在路上的数据怎么办。 取消是一个动作,而数据回传是持续的流;两者之间存在时间差,这个时间差里发生的事才是 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),然后在 aboutToDisappeartaskpool.cancel(this.task)。这段代码过不了编译,因为 execute() 返回的是 Promise,而 cancel() 的形参类型是 taskpool.Task。要能取消,必须显式 new taskpool.Task(...),把 Task 对象留下来。

把这三类方案里的错误集中起来,就得到了本文要逐条纠正的清单:

  • taskpool.execute() 的返回值是 Promise不可传给 cancel()
  • sendData / isCanceledtaskpool.Task静态方法,既不是 taskpool 上的方法,也不是注入进任务函数的回调;
  • sendData 只能在并发函数自己的执行线程内被调用,放进 setTimeoutPromise.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 / @TracemakeObserved),子线程修改共享对象的属性、主线程观测到变化并驱动刷新,是可以做到的。但请注意这条路径的实质:驱动刷新的仍然是主线程侧的观察者,不是子线程自己去碰渲染管线。它把「数据的跨线程共享」和「刷新的触发」这两件事分开处理了,铁律本身并没有被推翻。

二、三个会让方案直接跑不起来的 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 只能在并发函数自己的执行线程内调用,把它挪进 setTimeoutPromise.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 版本有关,较低版本对运行中任务的取消可能失败并抛出 BusinessError10200016 就属于这类「任务正在执行,无法取消」的情形)。所以 cancel() 一定要包 try / catch,而且永远不要把「取消成功」当成「不会再有回调」的保证

这就是为什么每个页面上都需要一个看起来有点多余的东西——代次(generation)校验。页面侧维护一个自增的 gen,每发起一次新任务就 ++gen 并记下本次的 myGen;所有异步回调(进度、结果、异常)进入时先比一次 myGen === this.gen,不相等就直接丢弃。再配合一个 alive 存活标记(aboutToAppeartrueaboutToDisappearfalse),就同时挡住了三种脏数据:

  • 页面已销毁后到达的进度与结果(alivefalse);
  • 上一次任务在取消后继续发出来的进度(gen 已递增);
  • 快速连点「开始」时,前一次任务迟到的结果覆盖后一次的结果(gen 已递增)。

需要强调的是,alivegen 校验不能省,哪怕已经正确调用了 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>,再用一个薄封装把「订阅 / 注销 / 同主题去重」收口到一处,避免 onoff 散落在各处。封装要守住三条:同一 scope + topic 重复订阅时先释放旧订阅;注销时用保存下来的 handler 引用做精确注销;页面销毁时按 scope 批量释放。

那 emitter 和 sendData 该怎么分工?我的建议很明确:同一个任务内部的进度与结果,优先用 sendData / onReceiveDataexecute() 的返回值,因为它们天然绑定任务归属,注销和去重都不需要你操心(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 里对 alivegen 的双重校验(取消后仍可能收到进度);finally 里按 Task 引用相等来判断是否清空(防止误清掉新任务);以及取消按钮里先 ++generationcancel 的顺序。

示例二:统一的 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)
  }
}

这里把「停止」和「销毁线程」拆成了两个按钮,就是为了体现前面的结论:长驻线程上的任务取消,应该优先走协作式停止(发指令 + 线程内检查标志位),而不是一上来就 terminateterminate 会让整个线程连同它维护的状态一起消失,重建成本高;只有在页面真正销毁、或者线程已经进入不可恢复状态时,才应该销毁它。

总结

把上面的内容压缩成一份可以贴在工位上的清单。

关于边界。 「子线程不能操作 UI」约束的是状态变量写入、渲染上下文的使用、以及编译期的可见性;它不约束计算本身。ArkTS 里没有、也不会有 setState 那种「调用入队、调度器兜底」的线程安全赋值,因为状态赋值的语义就是直接写字段,中间没有任何调度层。跨线程共享对象 + 主线程观察者的组合可以把边界往外推一点,但驱动刷新的仍然是主线程侧的观察者。

关于 API。 taskpool.execute() 返回的是 Promise,不是 Task;要能取消,必须先 new taskpool.Task(...) 并留住对象。sendDataisCanceled 都是 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 用起来其实比想象中干净得多。

Logo

讨论HarmonyOS开发技术,专注于API与组件、DevEco Studio、测试、元服务和应用上架分发等。

更多推荐