鸿蒙的 TaskPool 与 Worker:多线程并发该怎么选、怎么写

编者这周集中看了鸿蒙的多线程并发,重点落在 TaskPool 和 Worker 这两个工具上。看之前以为它们的关系是"一个新一个旧",或者"一个封装了另一个",看完才发现两者面向的维度根本不同,选错了不会编译报错,只会在运行时以各种难以理解的方式出问题:任务莫名其妙不执行了,数据传过去变成了空对象,闭包里的变量取不到值。这篇文章把机制、对比和踩坑整理一遍,给自己留个记录。

为什么需要把任务挪到别的线程

鸿蒙应用的默认执行环境是主线程,也就是宿主线程上加 UI 渲染、事件分发、生命周期回调全都跑在这里。一旦有个耗时的计算或者大量的数据解析占了主线程,界面就会卡住不动,因为负责刷新界面的活儿也排在同一个队列里等着。

解决方式很朴素,把耗时的部分挪到别的线程去算,算完了把结果交回主线程。鸿蒙为这件事提供了两套工具,TaskPool 和 Worker,两者的设计出发点并不一样。

TaskPool 的运作机制

TaskPool 提供的是任务队列式的多线程环境。宿主线程提交任务到队列,系统从工作线程里挑一个合适的来执行,执行完把结果返回给宿主线程。整个过程里开发者不需要关心线程的创建、复用和销毁。

官方文档里对它有几条描述值得单独拎出来。一是工作线程数量由设备的物理核数决定,内部自己管理具体数量;二是系统默认启动一个任务工作线程,任务多了会自动扩容;三是长时间没有任务分发时会缩容,减少工作线程数量。也就是说线程池的大小是动态的,跟着负载走。这一点直接决定了它的适用场景:任务量零散、时多时少、彼此独立的时候,让系统来调度比自己管线程划算。

TaskPool 的接口也围绕"任务"这个粒度展开。可以给任务设优先级,可以取消任务,可以把一组有关联的任务打包成任务组一次执行,也可以串行执行一批任务。这些都是任务级别的能力,不是线程级别的。

Worker 的运作机制

Worker 提供的是一个和宿主线程分离的独立执行环境。创建 Worker 的那个线程叫宿主线程,不一定是主线程,Worker 线程本身也支持再创建子 Worker。每个 Worker 子线程拥有独立的实例,包含独立的执行环境、对象和代码段。

这句话的份量在于最后半句。因为执行环境是独立的,所以每次启动一个 Worker 都会带来一定的内存开销,官方文档里明确写了需要限制 Worker 子线程的数量。这是它和 TaskPool 最根本的差别:TaskPool 的线程是共享复用的,Worker 的实例是各自独立的。

宿主线程和 Worker 之间通过消息传递机制通信,底层依赖序列化、引用传递或者转移所有权这几种机制来完成命令和数据的交互。Worker 的生命周期需要开发者自己管理,用完要主动终止。

两者的对比

官方文档给了一段相当明确的对比结论,这里摘其要点。TaskPool 的工作线程绑定系统的调度优先级,并且支持负载均衡即自动扩缩容;Worker 需要开发者自行创建,存在创建耗时,且不支持设置调度优先级。因此在性能方面使用 TaskPool 会优于 Worker,大多数场景推荐使用 TaskPool。

但"大多数场景"不等于"所有场景",下面几种情况必须用 Worker。

第一种是运行时间超过三分钟的任务。文档里特别注明,这三分钟不包含 Promise 和 async 与 await 这类异步调用的耗时,比如网络下载、文件读写这些 I/O 任务的等待时间不计入。举例来说,后台跑一个一小时的预测算法训练这类 CPU 密集型任务,就得用 Worker。

第二种是有关联的一系列同步任务。典型场景是需要创建句柄的场景,句柄每次创建都不同,需要永久保存下来并保证后续操作用的是同一个句柄,这种跨任务的持续状态 TaskPool 管不了。

反过来说,有几种情况用 TaskPool 更合适。需要设置优先级的任务,比如后台计算的直方图数据要用于前台界面显示、影响用户体验、需要高优先级处理。需要频繁取消的任务,比如大图浏览时缓存当前图片左右各两张,往一侧滑动时就要取消另一侧的缓存任务。大量或者调度点比较分散的任务,比如大型应用多个模块各自包含耗时任务,用 Worker 去做负载管理很不方便。

用一个表格概括就是:

判断依据用 TaskPool用 Worker
单次执行时长三分钟以内超过三分钟
是否需要设置优先级需要不支持
是否需要频繁取消需要不支持
是否需要长期持有状态不需要需要
任务数量与调度点多且分散少而集中

三分钟限制

TaskPool 的注意事项里有一条最容易被忽略。实现任务的函数在工作线程中的执行时长不能超过三分钟,否则如果因为任务逻辑导致阻塞、使任务无法完成,会导致该线程后续无法调度其他任务。当所有线程都被超时占用时,后续提交的任务将无法正常调度执行。

这段话描述的是一种连锁后果:一个写错的任务不是只让这一次调用失败,它会把整个线程池拖死。所以三分钟不是一个宽松的性能建议,而是一条必须守住的边界。

文档还进一步说明了这三分钟具体统计什么。它仅统计同步执行时长,不包含异步操作的等待时长。数据库的插入、删除、更新如果是异步操作,只计入 CPU 实际处理时长,比如 SQL 解析的时间,网络传输或磁盘 I/O 的等待时间不计入;如果是同步操作,整个操作时长包含 I/O 阻塞时间都要计。开发者可以通过 Task 的 ioDuration 和 cpuDuration 属性获取当前任务的异步 I/O 耗时和 CPU 耗时,用来判断自己有没有踩线。

如果确实需要长时间运行,TaskPool 提供了 LongTask(API 12 起)。文档里说它不设置执行时间上限,长时间运行不会触发超时异常,但不支持将同一任务多次执行,也不支持加入任务组。执行长时任务的线程会持续存在,直到任务完成并调用 terminateTask 后,该线程在空闲时才会被回收。这一点很关键,LongTask 的线程不会自动缩容回收,必须显式终止。

任务函数的写法

TaskPool 里实现任务的函数必须用 @Concurrent 装饰器标注,而且只支持在 .ets 文件中使用。这是导入和使用的基本样子:

import { taskpool } from '@kit.ArkTS';

@Concurrent
function computeSum(count: number): number {
  let sum: number = 0;
  for (let i = 0; i < count; i++) {
    sum += i;
  }
  return sum;
}

taskpool.execute(computeSum, 100).then((value: Object) => {
  console.info('任务结果: ' + value);
});

这里用的是 execute 直接传函数和参数的写法。文档里对它有明确说明:这种模式把待执行的函数放入内部任务队列,函数不会立即执行而是等待分发到工作线程,并且在当前执行模式下不支持取消任务。

要能取消、要能设优先级,就得先把任务构造出来:

import { taskpool } from '@kit.ArkTS';

let task: taskpool.Task = new taskpool.Task(computeSum, 100);
// execute 的第二个参数是优先级,默认是 taskpool.Priority.MEDIUM
taskpool.execute(task, taskpool.Priority.HIGH).then((value: Object) => {
  console.info('任务结果: ' + value);
});
// 需要时取消
taskpool.cancel(task);

文档里 execute 重载的说明是:将创建好的任务添加到内部任务队列,支持设置任务优先级和通过 cancel 取消任务;任务不能是任务组任务、串行队列任务或异步队列任务;长时任务只能调用一次,非长时任务可以多次调用执行。Task 的构造函数要求传入的函数必须用 @Concurrent 装饰。

如果需要一组有关联的任务,可以用 TaskGroup:

import { taskpool } from '@kit.ArkTS';

let taskGroup: taskpool.TaskGroup = new taskpool.TaskGroup();
taskGroup.addTask(computeSum, 100);
taskGroup.addTask(computeSum, 200);

taskpool.execute(taskGroup).then((values: Object[]) => {
  console.info('任务组结果: ' + JSON.stringify(values));
});

文档里说任务组一次执行一组任务,适用于执行一组有关联的任务。如果所有任务正常执行,异步执行完毕后返回所有任务结果的数组,数组中元素的顺序与 addTask 的顺序相同;如果任意任务失败则抛出对应异常,多个任务失败时抛出第一个失败任务的异常。任务组可以多次执行,但执行后不能新增任务。

任务和宿主线程之间怎么通信

有一种需求很常见:任务还在跑的时候,想让它把进度报回宿主线程,而不是等它全部做完。比如加载三十组图片,希望每加载完一组就在界面上更新一次。

TaskPool 为此提供了 sendData 和 onReceiveData 这一对接口。任务内部调用 taskpool.Task.sendData 把数据发出去,宿主线程侧通过任务的 onReceiveData 注册一个普通函数来接收:

import { taskpool } from '@kit.ArkTS';

function notice(data: number): void {
  console.info('子线程任务已执行完,共加载图片: ', data);
}

@Concurrent
function loadPictureSendData(count: number): string[] {
  let result: string[] = [];
  for (let index = 0; index < count; index++) {
    result.push(`item${index}`);
    // 即时通知宿主线程进度
    taskpool.Task.sendData(result.length);
  }
  return result;
}

let loadPictureTask: taskpool.Task = new taskpool.Task(loadPictureSendData, 30);
// 设置接收 Task 发送消息的方法
loadPictureTask.onReceiveData(notice);
taskpool.execute(loadPictureTask).then((res: Object) => {
  console.info('最终结果: ' + JSON.stringify(res));
});

注意 onReceiveData 注册的 notice 是个运行在宿主线程上的普通函数,它不需要 @Concurrent 装饰,因为它在宿主线程执行,不是任务的一部分。

Worker 的创建与通信

Worker 的第一步是把线程文件放对位置。文档里的要求是 Worker 线程文件需要放在 {moduleName}/src/main/ets/ 目录层级之下,否则不会被打包到应用中。手动创建的话一般是在 ets 目录下建一个 workers 文件夹存放 worker.ets,然后在 build-profile.json5 里配置:

"buildOption": {
  "sourceOption": {
    "workers": [
      './src/main/ets/workers/worker.ets'
    ]
  }
}

推荐的做法是用 DevEco Studio 自动创建,在模块目录下右键选择 New 再选 Worker,工具会把模板文件和上面的配置一起生成好,省掉手工配置容易漏字段的问题。

宿主线程这边构造 Worker 实例需要传线程文件的路径:

import { worker } from '@kit.ArkTS';

const workerInstance: worker.ThreadWorker = new worker.ThreadWorker('entry/ets/workers/worker.ets');

路径有三种写法。第一种是以 {moduleName}/ets/{relativePath} 的形式加载,relativePath 是相对于 {moduleName}/src/main/ets/ 的路径。第二种是以 @{moduleName}/ets/{relativePath} 的形式加载,带 @ 标识。第三种是相对路径形式,比如 …/…/workers/worker.ets,但这种写法仅支持包内加载,不支持跨包加载。文档的建议是加载 entry、feature 及 hsp 包的 Worker 线程文件时不建议用相对路径,推荐用第一种,不需要拼接路径。另外文件后缀 .ets 和 .ts 都可以省略。

Worker 线程文件内部的写法是这样的:

import { ErrorEvent, MessageEvents, ThreadWorkerGlobalScope, worker } from '@kit.ArkTS';

const workerPort: ThreadWorkerGlobalScope = worker.workerPort;

workerPort.onmessage = (e: MessageEvents) => {
  // 处理宿主线程发来的消息
  if (e.data === 'hello world') {
    workerPort.postMessage('success');
  }
};

注意这里和宿主线程侧的 API 名称上的区别。Worker 线程内部用的是 workerPort,它同时承担收发两个方向;宿主线程侧则是 Worker 实例上的 postMessage 和 onmessage,外加 onexit 用来感知 Worker 的退出。

import { worker } from '@kit.ArkTS';

async function postMessageTest(): Promise<void> {
  let ss: worker.ThreadWorker = new worker.ThreadWorker('entry/ets/workers/Worker.ets');
  let isTerminate: boolean = false;

  ss.onexit = () => {
    isTerminate = true;
  };

  // 接收 Worker 线程发送的消息
  ss.onmessage = (e) => {
    console.info('worker:: res is  ' + e.data);
  };

  // 给 Worker 线程发送消息
  ss.postMessage('hello world');

  // 用完主动终止,释放 Worker 实例的内存
  ss.terminate();
}

最后那句 terminate 是 Worker 与 TaskPool 在资源管理上最直观的差别。TaskPool 的线程由系统按负载扩缩容,开发者不用管;Worker 的实例创建出来就占着内存,不用了必须自己终止,否则既浪费内存又让线程一直挂着。

跨线程到底能传什么

这是编者这周踩得最深的一块,也是两套工具共同的约束。

TaskPool 的文档里写着,实现任务的函数入参需要满足序列化支持的类型。目前不支持使用 @State、@Prop、@Link 等装饰器修饰的复杂类型。这条限制很明确:状态管理的那些响应式对象不能跨线程传,因为它们身上绑定了 UI 框架的更新逻辑,既无法序列化也没有意义。

另外从 API 11 开始,跨并发实例传递带方法的实例对象时,该类必须使用 @Sendable 装饰器标注,且仅支持在 .ets 文件中使用。如果不想用 @Sendable,文档给的建议是考虑改用 Worker,因为 Worker 支持同步调用宿主线程的接口。

ArrayBuffer 的传递行为需要特别留意。在 TaskPool 中 ArrayBuffer 参数默认是转移而不是复制,需要设置转移列表时可以通过 setTransferList 接口设置。转移的意思是 buffer 的控制权交给了工作线程,传输之后当前这个 ArrayBuffer 就失效了。如果确实需要多次调用一个以 ArrayBuffer 作参数的 task,就得通过 setCloneList 把传输行为改成拷贝传递,避免对原有对象产生影响。文档里还提醒 setCloneList 需要搭配 @Sendable 装饰器使用,否则会抛异常。

Worker 这边的通信对象类型限制,官方文档单独列了一份文档讲线程间通信对象,涵盖普通对象、容器类对象、ArrayBuffer、SharedArrayBuffer、Transferable 对象和 Sendable 对象,各自的传递语义不一样。篇幅所限这里不展开,但结论是一样的:跨线程传的函数入参要满足序列化支持的类型,带方法的对象要么用 @Sendable 标注,要么走 Worker 的同步调用。

小结

如果只记一句话,编者会这样总结:TaskPool 是在线程池里提交一个个任务,任务结束线程就还回去,适合短平快、数量多、可能要取消或排优先级的场景;Worker 是开一个独立的常驻执行环境,适合跑得久、需要长期持有状态、或者需要同步调用宿主线程接口的场景。

而两者共通的约束只有一条,跨线程的数据必须能序列化,带方法的对象要么用 @Sendable 标注,要么绕开这个需求。理解了这条,前面那些数据传不过去、传过去变成空对象的怪现象,基本都能解释清楚。

Logo

作为“人工智能6S店”的官方数字引擎,为AI开发者与企业提供一个覆盖软硬件全栈、一站式门户。

更多推荐