鸿蒙并发模型 TaskPool 与 Worker

本文讲清 HarmonyOS NEXT 的并发模型:ArkTS 主线程如何被耗时任务拖垮,TaskPool 与 Worker 各自适合什么场景。覆盖 @Concurrent 装饰器、Task 对象、TaskGroup、优先级与取消、Worker 的创建与通信、Sendable 对象与 collections 共享容器,并用对照表给出选型依据与高频坑清单。

开篇:为什么 UI 线程不能做重活

ArkTS 的运行时是单线程事件循环模型:所有 UI 绘制、手势响应、状态刷新都跑在主线程上。一旦主线程被一段长任务占住,事件循环就无法把控制权交还给渲染管线,用户看到的就是白屏、掉帧、点击无响应。

这个问题在鸿蒙上比在 Web 上更尖锐,因为 ArkUI 的声明式渲染把"状态变化到界面更新"压缩在同一个线程里,主线程卡顿会同时冻结数据流与视图流。因此,把耗时计算移出主线程不是优化项,而是必须遵守的底线。

本文覆盖两条并发路径:轻量的 TaskPool 与常驻的 Worker,以及支撑它们传递数据的 Sendable 机制。示例基于 API 12 及以上,导入路径使用 @kit.ArkTS。如果你对 ArkTS 的类型约束还不熟,建议先读 ArkTS 语言基础与 TypeScript 的差异 。

一、ArkTS 的线程模型

1.1 主线程与事件循环

主线程维护一个任务队列,按顺序执行微任务与宏任务。UI 渲染请求、状态变更回调、定时器、网络回调都会进入这个队列。只要某个任务执行时间过长,后续任务全部延后。

1.2 什么算耗时任务

判断标准不是"代码长不长",而是"单次执行是否超过一帧的预算"。在 60 帧目标下,一帧只有约 16.7 毫秒,留给业务逻辑的空间通常不到 5 毫秒。

任务类型典型耗时是否必须移出主线程
简单状态赋值微秒级否
JSON 解析 1MB十毫秒级视情况
图片解码与压缩百毫秒级是
大数组排序 10 万条百毫秒级是
加密与哈希计算百毫秒级是
数据库批量写入秒级是

二、TaskPool 详解

TaskPool 是鸿蒙提供的托管任务池,由系统统一调度线程数量,开发者只负责提交任务。它适合"一次性、可并行、无状态"的计算。

2.1 @Concurrent 与 execute

并发函数必须用 @Concurrent 装饰,且必须是顶层函数或类的静态方法。

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

@Concurrent
function sumRange(start: number, end: number): number {
  let total = 0;
  for (let i = start; i <= end; i++) {
    total += i;
  }
  return total;
}

export async function computeSum(): Promise<number> {
  const result = await taskpool.execute(sumRange, 1, 1000000);
  return result as number;
}

taskpool.execute 有两种重载:直接传函数与参数,或传入一个 taskpool.Task 对象。前者写法简洁,后者能控制优先级、任务组与取消。

2.2 Task 对象与参数传递

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

@Concurrent
function compressImage(pixelData: ArrayBuffer, quality: number): ArrayBuffer {
  // 在子线程中执行压缩,避免阻塞 UI
  return pixelData;
}

export async function compress(pixelData: ArrayBuffer, quality: number): Promise<ArrayBuffer> {
  const task = new taskpool.Task(compressImage, pixelData, quality);
  task.setTransferList([pixelData]);
  const result = await taskpool.execute(task);
  return result as ArrayBuffer;
}

setTransferList 是关键优化:列进 transfer list 的 ArrayBuffer 会以转移所有权的方式传递,避免整块内存拷贝。对于几十 MB 的像素数据,这一步能省下可观的耗时与峰值内存。

2.3 任务组与优先级

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

@Concurrent
function processChunk(chunk: string): number {
  return chunk.length;
}

export async function processAll(chunks: string[]): Promise<number[]> {
  const group = new taskpool.TaskGroup();
  for (const chunk of chunks) {
    const task = new taskpool.Task(processChunk, chunk);
    task.setPriority(taskpool.Priority.HIGH);
    group.addTask(task);
  }
  const results = await taskpool.execute(group);
  return results as number[];
}

taskpool.Priority 有 HIGH、MEDIUM、LOW 三档,默认 MEDIUM。优先级只影响同池内的排队顺序,不要指望用高优先级来"抢占"已经在跑的任务。

2.4 取消与超时

Task 对象可以调用 cancel() 取消尚未执行的任务,但如果任务已经开始执行,取消不会中断正在运行的代码。因此长任务应该在函数内部主动检查一个可传递的中断标志,而不是依赖外部强杀。

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

@Concurrent
function pollUntil(deadline: number): string {
  let step = 0;
  while (Date.now() < deadline) {
    // 主动检查截止时间,给取消留出退出点
    step++;
  }
  return `steps: ${step}`;
}

export function scheduleDelayed(): taskpool.Task {
  const task = new taskpool.Task(pollUntil, Date.now() + 5000);
  taskpool.executeDelayed(2000, task).then((value: Object) => {
    console.info(value as string);
  }).catch((err: Error) => {
    console.error(`task failed: ${err.message}`);
  });
  return task;
}

executeDelayed(delay, task) 的第一个参数是延迟毫秒数,用于给批量任务做节流。拿到返回的 Task 后可以随时 taskpool.cancel(task),但如前所述,它只能把还没开始的任务从队列里摘掉。

2.5 任务池的容量与排队

TaskPool 的线程数量由系统根据设备核心数与当前负载动态决定,开发者无法直接指定。这意味着大量任务同时提交时会出现排队,实测中如果同时提交几百个微任务,后面的任务可能等待数秒才被执行。

提交规模现象建议
几个任务立即执行直接提交
几十个任务轻微排队用 TaskGroup 批量提交
几百个任务明显排队拆成批次,用 executeDelayed 错峰
上千个任务调度开销显著改用手动分片,减少任务数

三、Worker 详解

Worker 是常驻线程,适合需要保持状态、反复通信的场景,例如持续接收 WebSocket 数据并在后台做聚合。

3.1 创建与通信

Worker 需要一个独立的入口文件,通过字符串路径引用。

// entry/src/main/ets/workers/CalcWorker.ets
import { worker, ThreadWorkerGlobalScope, MessageEvents } from '@kit.ArkTS';

const workerPort: ThreadWorkerGlobalScope = worker.workerPort;

workerPort.onmessage = (e: MessageEvents) => {
  const input = e.data as number[];
  let total = 0;
  for (const n of input) {
    total += n;
  }
  workerPort.postMessage(total);
};
// 主线程侧
import { worker } from '@kit.ArkTS';

export class CalcClient {
  private instance: worker.ThreadWorker | null = null;

  start(onResult: (total: number) => void): void {
    this.instance = new worker.ThreadWorker('entry/ets/workers/CalcWorker.ets');
    this.instance.onmessage = (e: MessageEvents) => {
      onResult(e.data as number);
    };
    this.instance.onerror = (e: ErrorEvent) => {
      console.error(`worker error: ${e.message}`);
    };
  }

  send(data: number[]): void {
    this.instance?.postMessage(data);
  }

  stop(): void {
    this.instance?.terminate();
    this.instance = null;
  }
}

3.2 Worker 的生命周期管理

Worker 不会自动回收,必须显式 terminate()。页面 onPageHide 或组件 aboutToDisappear 时如果没有终止 Worker,线程会一直存活并持有上下文引用,造成内存泄漏与耗电。这一点在页面频繁进出时尤其明显。

四、Sendable 对象与共享内存

TaskPool 与 Worker 默认走序列化传值,深拷贝的代价在大对象上不可忽视。Sendable 机制允许对象在并发实例间以引用方式共享,从而绕开拷贝。

4.1 @Sendable 装饰器

@Sendable
export class Counter {
  count: number = 0;

  increment(): void {
    this.count += 1;
  }
}

@Concurrent
function bump(counter: Counter, times: number): number {
  for (let i = 0; i < times; i++) {
    counter.increment();
  }
  return counter.count;
}

被 @Sendable 装饰的类有严格限制:不能有方法外的成员、不能继承非 Sendable 类、属性类型必须是 Sendable 支持的简单类型或其他 Sendable 对象。

4.2 collections 的 Sendable 容器

@kit.ArkTS 提供的 collections 模块包含 Sendable 版本的 Array、Map、Set,可以在并发实例间安全共享。

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

@Concurrent
function pushAll(target: collections.Array<number>, items: number[]): void {
  for (const item of items) {
    target.push(item);
  }
}

普通 Array 与 collections.Array 不能混用,前者传进并发函数会被序列化拷贝,后者才是真正的共享引用。

4.3 共享数据的同步约束

共享内存省掉了拷贝,但也把竞态问题带了回来。ArkTS 目前没有提供语言级的互斥锁,Sendable 对象的并发读写只能靠"单写者"约定来规避:同一时刻只允许一个并发实例修改同一对象,其他实例只读。如果做不到,就应该退回序列化传值,用拷贝换安全。

数据形态传递方式是否拷贝并发安全
简单类型参数序列化是天然安全
普通对象与数组序列化是天然安全
@Sendable 对象引用共享否需约定单写者
collections.Array引用共享否需约定单写者
ArrayBuffer 加 transferList所有权转移否转移后原引用失效

实践中更稳妥的做法是:共享容器只用于"主线程写、子线程读"的单向场景,或者让每个并发任务各自持有一份局部数据,最后由主线程合并。

五、TaskPool 与 Worker 选型

维度TaskPoolWorker
线程模型系统托管线程池开发者创建的常驻线程
启动开销低,任务粒度较高,需创建入口文件
状态保持无状态,任务间不共享可保持内部状态
通信方式返回值与 PromisepostMessage 双向
任务取消支持,仅未执行时可取消需自行实现中断协议
并行数量由系统调度每个实例一个线程
适用场景一次性计算、并行批处理长连接处理、持续聚合
典型例子图片压缩、排序、哈希WebSocket 数据聚合

选型原则:能用 TaskPool 就不要用 Worker。TaskPool 的线程由系统托管,能自动适配设备核心数,而 Worker 需要自己控制实例数量,创建过多反而会因线程切换拖慢整体。只有当任务需要长期保持状态、或需要与主线程持续双向通信时,才选择 Worker。

六、实战:图片压缩流水线

把上面的能力拼成一个可用的批处理流程,思路是把大图切成多个分片并行压缩,最后在主线程合并。

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

@Concurrent
function compressChunk(pixels: ArrayBuffer, quality: number): ArrayBuffer {
  // 实际项目在此调用图像处理库
  return pixels;
}

export async function compressBatch(
  chunks: ArrayBuffer[],
  quality: number
): Promise<ArrayBuffer[]> {
  const group = new taskpool.TaskGroup();
  for (const chunk of chunks) {
    const task = new taskpool.Task(compressChunk, chunk, quality);
    task.setTransferList([chunk]);
    task.setPriority(taskpool.Priority.MEDIUM);
    group.addTask(task);
  }
  const results = await taskpool.execute(group);
  return results as ArrayBuffer[];
}

主线程侧的调用方式很直接:把耗时逻辑交给并发函数,自己在 await 之后更新状态。

@Entry
@Component
struct CompressPage {
  @State progress: number = 0;

  private async splitPixels(): Promise<ArrayBuffer[]> {
    // 从相册读取并按分片切分像素数据
    return [];
  }

  async doCompress(): Promise<void> {
    const chunks = await this.splitPixels();
    this.progress = 30;
    const results = await compressBatch(chunks, 80);
    this.progress = 100;
    console.info(`compressed ${results.length} chunks`);
  }

  build() {
    Column({ space: 16 }) {
      Progress({ value: this.progress, total: 100 })
        .width('80%')
      Button('开始压缩')
        .onClick(() => {
          this.doCompress();
        })
    }
    .width('100%')
    .padding(24)
  }
}

这套写法有两个收益:主线程全程只做 Promise 等待,帧率不受影响;setTransferList 让大缓冲区不产生额外拷贝。代价是并发度交给系统决定,极端情况下任务排队时间不可控,对时延敏感的场景需要配合 executeDelayed 做节流。

需要注意的是,并发函数内部不能直接更新 UI 状态。它只能返回结果,由主线程在 await 之后赋值给状态变量。这条约束和鸿蒙应用与页面生命周期里讲的线程边界是同一件事。

并发框架本身也会消耗内存与线程资源,如果任务粒度太细(例如每个任务只算几十个数的和),调度开销会超过计算本身。判断标准是单任务耗时是否显著大于任务提交开销,通常以 1 毫秒为界。相关的量化方法可以参考 鸿蒙原生应用性能优化与调试 。

在数据结构的选择上,网络层返回的大 JSON 往往需要先解析再处理,解析与处理都应放在子线程,主线程只接收最终结果。数据层的整体设计见 鸿蒙网络请求与数据持久化 。如果你熟悉 Dart 的 Isolate 模型,会发现两者的取舍逻辑高度相似。

七、常见坑清单

  • @Concurrent 函数里引用了闭包变量或类成员,编译期直接报错。
  • 传入的参数包含不可序列化的对象(如 context、函数),运行时抛异常。
  • @Concurrent 函数内部直接修改状态变量,编译不通过或修改无效。
  • Worker 路径字符串写错,运行时找不到入口文件。
  • Worker 忘记 terminate(),页面反复进出后线程堆积。
  • 任务粒度过细,调度开销大于计算收益,反而更慢。
  • 误用普通 Array 期待共享,实际发生的是深拷贝。
  • @Sendable 类里放了非 Sendable 类型的字段,编译报错。
  • 依赖 cancel() 中断已在执行的长任务,实际无法生效。
  • 在 Worker 中访问主线程的 UI 对象,触发运行时错误。

小结

鸿蒙的并发模型只给了两条路,选择标准也很清晰:一次性、无状态、可并行的计算交给 TaskPool,需要常驻与持续通信的场景才用 Worker。真正决定性能的不是用了哪个 API,而是有没有把任务粒度、数据传递方式和生命周期管理想清楚。记住三条底线:并发函数必须是顶层或静态函数且不捕获上下文,大缓冲区用 setTransferList 或 Sendable 避免拷贝,Worker 用完必须 terminate()。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「鸿蒙开发」更多文章

  1. 鸿蒙 ohpm 包管理与 Hypium 测试框架
  2. ArkUI 动画体系与手势交互
  3. 鸿蒙应用安全:权限模型与 HUKS 密钥管理