引言

如果你开发过一个计算密集型的鸿蒙应用——比如图片处理、数据分析、密码学运算——你大概率遇到过这样的场景:用户点击"开始处理"按钮,然后整个界面就卡死了。按钮没有反馈,动画停止播放,滑动毫无反应,用户体验极其糟糕。原因很简单:JavaScript 本质上是单线程的,所有的 UI 渲染和业务逻辑都跑在同一个主线程上。当一个耗时计算占用了主线程,UI 就无法响应。

前端开发者对此并不陌生。在 Web 端,我们有 Web Workers 来解决这个问题;在 Node.js 中,有 worker_threads 模块。HarmonyOS 同样提供了完整的多线程并发方案:@ohos.worker@ohos.taskpool 两个模块。前者提供了 ThreadWorker 类——可以创建独立的线程来执行耗时任务,通过消息机制与主线程通信,并且支持随时终止线程。后者 TaskPool 则是一个更轻量的并发任务池,适合独立的短任务。

本文将通过构建一个"并发计算实验室",深入讲解 @ohos.worker 的核心 API:ThreadWorker 创建、postMessage/onmessage 双向通信、onerror 错误处理、terminate() 线程终止。Demo 实现了 Worker 线程完成质数计算和斐波那契数列计算,并与主线程的执行效果进行对比——让读者直观感受多线程的价值。

读完本文你将能够:

  • 使用 worker.ThreadWorker 创建独立线程并加载 worker 脚本
  • 使用 postMessage() / onmessage 在主线程与 Worker 之间双向通信
  • 在 Worker 中执行耗时计算并实时回传进度
  • 处理 Worker 错误和生命周期(onerrorterminate()
  • 理解 Worker 线程与主线程的性能差异
  • 了解 @ohos.taskpool 的轻量并发方案

为什么需要多线程:主线程的困境

在深入了解 Worker API 之前,先理解一下问题的本质。

ArkUI 应用的架构可以用一句话概括:一个主线程,负责所有 UI 渲染和事件响应。当你在 ArkTS 代码中执行一个耗时计算时,这个计算会阻塞主线程的事件循环。事件循环被阻塞意味着:

  1. UI 无法更新@State 变量虽然已经修改,但框架没有机会重新渲染。
  2. 用户交互无响应:点击、滑动、返回等手势事件被堆积在事件队列中,没有机会被处理。
  3. 系统 ANR 风险:如果主线程阻塞超过 5 秒,系统可能会弹出"应用无响应"对话框。

下面这个简单的例子就能说明问题:

// 在主线程中计算 100 万以内的质数 —— UI 冻结!
@Entry
@Component
struct BadExample {
  @State result: string = '计算中...';

  aboutToAppear(): void {
    let primes: number[] = [];
    for (let i = 2; i <= 1000000; i++) {
      // 判断 i 是否为质数...
      if (isPrime(i)) primes.push(i);
    }
    this.result = '找到 ' + primes.length + ' 个质数';
  }
}

这段代码在 aboutToAppear 中执行了耗时计算。页面会一直显示白屏,直到计算完成——用户看到的是应用"卡住了"。

正确的做法是将耗时的质数计算放到 Worker 线程中执行,主线程只负责接收结果和更新 UI。

@ohos.worker 模块概述

两个并发模块的选择

HarmonyOS 提供了两个并发模块:

特性 @ohos.worker @ohos.taskpool
线程模型 独立线程(ThreadWorker) 系统管理的任务池
生命周期 手动创建和销毁 系统自动管理
通信方式 postMessage/onmessage Task 参数和返回值
进度报告 支持中间进度回传 不支持(仅最终结果)
取消支持 terminate() 强制终止 taskpool.cancel()
适用场景 长时计算、持续通信 独立短任务、并行处理
脚本文件 需要单独的 .ets 文件 不需要(函数即可)

简单原则: 如果你需要实时进度报告、需要在计算过程中与主线程交互、或者计算是持续的(如 WebSocket 心跳),用 Worker。如果只是执行一个独立的计算密集型任务(如图片压缩),用 TaskPool 更简单。

本文 Demo 选择 Worker,因为它能展示更完整的多线程编程模型——包括进度报告、线程生命周期管理、以及错误处理。

Worker 文件结构

Worker 需要一个独立的 .ets 文件作为线程执行脚本。这个文件运行在 Worker 线程的上下文中,不能访问主线程的 UI 组件,但可以使用所有 ArkTS 标准库。

典型的 Worker 文件结构:

// entry/src/main/ets/workers/MyWorker.ets
import { worker } from '@kit.ArkTS';

const workerPort = worker.workerPort;

// 接收主线程消息
workerPort.onmessage = (e: MessageEvents) => {
  let data = e.data;
  // 执行计算...
  
  // 回传结果
  workerPort.postMessage({ result: 'done' });
};

worker.workerPort 是 Worker 线程的全局通信端口,它是 Worker 侧与主线程通信的唯一桥梁。

核心 API 逐项解析

导入模块

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

Worker API 归属于 @kit.ArkTS 套件,与 UI 框架同一来源。这很合理——Worker 虽然不直接渲染 UI,但它在 ArkTS 应用的运行时环境中运行,需要框架提供线程管理能力。

ThreadWorker(path: string)

创建 Worker 实例并加载指定路径的 worker 脚本:

private workerInstance: worker.ThreadWorker | null = null;

// 创建 Worker
this.workerInstance = new worker.ThreadWorker('entry/ets/workers/ComputeWorker.ets');

路径规则:

  • 格式为 {moduleName}/ets/{relativePath}
  • 对于 entry 模块的 worker 脚本,路径以 entry/ets/workers/ 开头。
  • Worker 脚本必须在编译时可见(放在 src/main/ets/ 目录下)。

创建 ThreadWorker 实例后,底层会启动一个新的系统线程,加载并执行指定的 worker 脚本。这个线程拥有自己独立的 JS 运行时,与主线程完全隔离——全局变量、模块缓存、内存空间都是独立的。

postMessage(msg: Object)

向对方线程发送消息。主线程调用 workerInstance.postMessage(),Worker 通过 workerPort.onmessage 接收;Worker 调用 workerPort.postMessage(),主线程通过 workerInstance.onmessage 接收:

// 主线程 → Worker
this.workerInstance.postMessage({
  command: 'compute',
  limit: 300000,
  taskType: 'primes'
});

// Worker → 主线程(在 worker 脚本中)
workerPort.postMessage({
  type: 'progress',
  progress: 45,
  current: 135000
});

消息内容可以是任意可序列化的 JavaScript 对象。这意味着你可以传递数字、字符串、数组、普通对象——但不能传递函数、DOM 节点、或 ArkUI 组件引用。

实践中,建议设计一个简单的消息协议——每个消息包含 type 字段标识消息类型,便于接收方分发处理:

// 消息协议设计
interface WorkerMessage {
  type: 'start' | 'progress' | 'complete' | 'error';
  // 其他字段随 type 变化
}

onmessage 和 onerror

主线程通过 onmessage 接收 Worker 发来的消息,通过 onerror 处理 Worker 中的未捕获异常:

this.workerInstance.onmessage = (e: MessageEvents) => {
  let msg = e.data;
  if (msg.type === 'progress') {
    this.progress = msg.progress;
    this.currentPos = msg.current;
  } else if (msg.type === 'complete') {
    this.result = { count: msg.count, elapsed: msg.elapsed };
  }
};

this.workerInstance.onerror = (err: ErrorEvent) => {
  this.statusText = 'Worker 错误: ' + err.message;
  this.terminateWorker();
};

onerror 是 Worker 异常的最后一道防线。Worker 脚本中未被 try-catch 捕获的异常最终会触发 onerror 回调。收到错误后,建议调用 terminate() 清理线程资源。

terminate()

终止 Worker 线程并释放系统资源:

terminateWorker(): void {
  if (this.workerInstance) {
    this.workerInstance.terminate();
    this.workerInstance = null;
  }
}

调用 terminate() 后:

  • Worker 线程立即停止执行(强制终止,不等待当前任务完成)。
  • 与 Worker 的通信通道关闭。
  • 线程占用的系统资源(内存、CPU 时间片)被释放。
  • 后续不能再通过该实例 postMessage()——需要重新创建。

在组件生命周期中,务必在 aboutToDisappear() 中清理 Worker:

aboutToDisappear(): void {
  this.terminateWorker();
}

如果不清理,Worker 线程可能会在页面销毁后继续运行——造成内存泄漏和 CPU 资源浪费。
在这里插入图片描述
在这里插入图片描述

Demo 设计:并发计算实验室

本文 Demo 实现了一个完整的"并发计算实验室",包含以下功能:

页面结构

Column(根容器)
├── Header(深色标题栏:"并发计算实验室" + @ohos.worker 标签)
├── Scroll
│   └── Column
│       ├── Worker 状态卡(运行中/空闲 + 已连接/未连接 + 进度百分比)
│       ├── 任务配置区
│       │   ├── 计算类型选择(质数计算 / 斐波那契)
│       │   ├── 计算量滑块(5万 ~ 100万)
│       │   └── 启动/取消按钮
│       ├── 计算进度区
│       │   ├── Progress 进度条(实时更新)
│       │   ├── 当前处理位置 + 已找到数量
│       │   └── 计算结果(总数 + 耗时 + 最后5项)
│       ├── 主线程 vs Worker 对比区
│       │   ├── 说明文字(告诫主线程冻结风险)
│       │   ├── "主线程同步计算"按钮
│       │   └── 耗时对比 + 加速比
│       ├── Worker 多线程优势(4条编号说明)
│       ├── 运行日志(最多10条,标注来源:Worker/Main)
│       └── API 参考区
└── 根容器结束

5 个交互点

  1. Worker 线程计算:选择计算类型(质数/斐波那契)→ 调整计算量滑块 → 点击"Worker 线程计算"→ 进度条实时更新(Worker 每 1% 进度发一次消息)→ 计算完成显示总数和耗时。
  2. 取消计算:Worker 计算中 → 点击"取消"按钮 → 调用 terminate() → 线程终止 → 状态恢复为"已取消计算"。
  3. 主线程同步计算:点击"主线程同步计算"按钮 → 主线程直接执行相同的计算 → 界面冻结 2-5 秒 → 计算完成后显示耗时 → 自动计算 Worker 相比主线程的加速比。
  4. 调整计算量:拖动滑块改变计算量(5万/10万/…/100万)→ 标签实时显示当前选择 → 重新开始计算验证不同规模下的性能表现。
  5. 运行日志追踪:每次 Worker 创建、计算开始、进度更新、完成、取消都自动记录日志 → 标注来源(Worker 蓝色 / Main 橙色)→ 最多保留 10 条。

核心实现

Worker 脚本:质数计算与斐波那契

ComputeWorker.ets 是整个 Demo 的核心——它运行在独立线程中,负责执行耗时计算:

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

const workerPort = worker.workerPort;

// 质数判断
function isPrime(n: number): boolean {
  if (n < 2) return false;
  if (n === 2) return true;
  if (n % 2 === 0) return false;
  let sqrt = Math.sqrt(n);
  for (let i = 3; i <= sqrt; i += 2) {
    if (n % i === 0) return false;
  }
  return true;
}

// 查找质数——在计算过程中定期报告进度
function findPrimes(limit: number): number[] {
  let primes: number[] = [];
  let lastReport = 0;
  for (let i = 2; i <= limit; i++) {
    if (isPrime(i)) {
      primes.push(i);
    }
    let progress = Math.floor((i / limit) * 100);
    if (progress > lastReport) {
      lastReport = progress;
      workerPort.postMessage({
        type: 'progress',
        progress: progress,
        current: i,
        found: primes.length
      });
    }
  }
  return primes;
}

// 接收主线程消息,开始计算
workerPort.onmessage = (e: MessageEvents) => {
  let data = e.data;
  if (data.command === 'compute') {
    workerPort.postMessage({ type: 'start', limit: data.limit });

    let startTime = Date.now();
    let result = data.taskType === 'primes'
      ? findPrimes(data.limit)
      : computeFibonacci(data.limit);
    let elapsed = Date.now() - startTime;

    workerPort.postMessage({
      type: 'complete',
      count: result.length,
      elapsed: elapsed,
      sample: result.slice(-5)  // 最后5个结果作为样本
    });
  }
};

关键设计点:

  • 批量进度报告progress > lastReport 保证每 1% 的进度只发一次消息,避免消息风暴(100 万次迭代发 100 万条消息会直接瘫痪主线程)。
  • 时间统计Date.now() 在 Worker 内部计时——这测量的是纯计算时间,不包含消息传递开销。这个时间比主线程的耗时更有参考价值。
  • 结果采样:只回传最后 5 个结果作为样本(result.slice(-5)),而不是将全部结果(可能有数万个)传回主线程,减少消息序列化开销。

主线程:Worker 生命周期管理

主页面中,startWorker() 方法负责完整的 Worker 生命周期:

startWorker(): void {
  // 重置状态
  this.progress = 0;
  this.currentPos = 0;
  this.foundCount = 0;
  this.isRunning = true;
  this.statusText = 'Worker 线程启动中...';

  try {
    // 创建 Worker 实例
    this.workerInstance = new worker.ThreadWorker(
      'entry/ets/workers/ComputeWorker.ets'
    );
    this.addLog('Worker', '线程创建成功');

    // 注册消息回调
    this.workerInstance.onmessage = (e: MessageEvents) => {
      let msg = e.data;
      if (msg.type === 'start') {
        this.statusText = 'Worker 正在计算...';
      } else if (msg.type === 'progress') {
        this.progress = msg.progress;
        this.currentPos = msg.current;
        this.foundCount = msg.found;
      } else if (msg.type === 'complete') {
        this.result = { count: msg.count, elapsed: msg.elapsed, lastFive: msg.sample };
        this.progress = 100;
        this.isRunning = false;
        this.terminateWorker();
      }
    };

    // 错误处理
    this.workerInstance.onerror = (err: ErrorEvent) => {
      this.isRunning = false;
      this.statusText = 'Worker 错误: ' + err.message;
      this.terminateWorker();
    };

    // 发送计算命令
    this.workerInstance.postMessage({
      command: 'compute',
      limit: this.computeLimit,
      taskType: this.taskType
    });
  } catch (e) {
    this.isRunning = false;
    this.statusText = 'Worker 创建失败';
  }
}

生命周期流程:

创建 ThreadWorker → 加载脚本 → postMessage(任务参数)
  → onmessage 接收 start 消息
  → onmessage 接收多次 progress 消息(更新 UI)
  → onmessage 接收 complete 消息 → terminate() 清理线程

错误路径:

创建失败(catch)→ 显示错误状态
运行时错误(onerror)→ 显示错误信息 → terminate() 清理
用户取消(cancelCompute)→ terminate() 强制终止
页面销毁(aboutToDisappear)→ terminate() 防止泄漏

主线程对比:让用户直观感受差异

Demo 中有一个特别的设计——"主线程同步计算"按钮。它直接在 UI 线程执行相同的计算任务:

runOnMainThread(): void {
  this.mainThreadResult = new ComputeResult();
  this.mainThreadRunning = true;
  this.addLog('Main', '主线程开始计算 (上限:' + this.computeLimit + ')');

  let startTime = Date.now();
  let result: number[] = [];

  // 在主线程中直接执行质数计算——UI 将完全冻结!
  for (let i = 2; i <= this.computeLimit; i++) {
    let prime = true;
    if (i < 2) prime = false;
    if (i === 2) prime = true;
    else if (i % 2 === 0) prime = false;
    else {
      let sqrt = Math.sqrt(i);
      for (let j = 3; j <= sqrt; j += 2) {
        if (i % j === 0) { prime = false; break; }
      }
    }
    if (prime) result.push(i);
  }

  let elapsed = Date.now() - startTime;
  this.mainThreadResult = {
    count: result.length,
    elapsed: elapsed,
    lastFive: result.slice(-5)
  };
  this.mainThreadRunning = false;
  this.mainThreadWarning = '主线程在 ' + this.formatTime(elapsed) + ' 内无响应';

  // 计算加速比
  if (this.result.count > 0) {
    this.speedUp = parseFloat((elapsed / this.result.elapsed).toFixed(1));
  }
}

这个对比有几个有趣的现象:

  1. 按钮在计算期间无法点击:因为主线程完全被计算任务占用,事件循环被阻塞。
  2. 进度条不会更新:在计算过程中,Progress 组件的值不会变化——虽然你可以把进度更新代码插入到循环中,但没有事件循环来处理渲染。
  3. 耗时可能相似甚至更快:Worker 有线程创建和消息序列化的开销,对于较小的计算量(5 万以内),主线程可能更快。但对于大计算量(30 万以上),Worker 的非阻塞优势就显现出来了——不是更快,而是不卡

这是理解多线程价值的关键:Worker 的目标不是让计算变快,而是让 UI 在计算的同时保持响应。

消息协议设计

Demo 中设计了一套清晰的消息协议:

主线程 → Worker:
{ command: 'compute', limit: number, taskType: 'primes' | 'fib' }

Worker → 主线程:
// 任务开始
{ type: 'start', limit: number, taskType: string }

// 进度更新(多次)
{ type: 'progress', progress: number, current: number, found: number }

// 计算完成
{ type: 'complete', count: number, elapsed: number, sample: number[] }

这种 type 字段的消息路由模式是 Worker 通信的最佳实践。接收方通过 type 快速判断消息种类,避免在一个回调中堆砌复杂的条件逻辑。

@ohos.taskpool 简介:更轻量的并发方案

虽然本文 Demo 使用 @ohos.worker,但在实际项目中,很多并发场景使用 @ohos.taskpool 会更简单。

TaskPool 不需要单独的 worker 脚本文件——你只需要用 @Concurrent 装饰器标记一个函数,然后将它封装为 Task 提交给任务池执行:

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

@Concurrent
function compressImage(data: ArrayBuffer): ArrayBuffer {
  // 在独立线程中执行图片压缩
  let compressed = heavyCompress(data);
  return compressed;
}

// 提交任务
let task = new taskpool.Task(compressImage, imageData);
taskpool.execute(task).then((result: ArrayBuffer) => {
  // 主线程接收压缩结果
  this.compressedImage = result;
}).catch((err: Error) => {
  console.error('压缩失败: ' + err.message);
});

TaskPool 的优势:

  • 无需单独文件:函数直接定义在组件文件中,用 @Concurrent 装饰即可。
  • 自动线程管理:系统自动分配线程、排队任务、回收空闲线程。
  • Promise 风格.then() / .catch() 接收结果,代码更简洁。
  • 支持取消taskpool.cancel(task) 取消排队中或运行中的任务。

TaskPool 的局限:

  • 不能报告中间进度:函数执行完成前不会回传任何消息。
  • 不保证线程独占:多个 Task 可能共享同一个物理线程(由系统调度)。
  • 适合短任务:不适合需要持续运行的后台任务(如 WebSocket 心跳)。

选择建议:如果一个任务可以在几秒内完成且不需要中间进度报告,用 TaskPool。如果需要长时间运行、需要实时进度、或者需要在任务运行期间与主线程多次交互,用 Worker

实际应用场景

场景一:图片批量处理

用户在相册中选择了 50 张照片,需要压缩为缩略图。在 Worker 中逐张处理,每处理完一张就发消息更新进度条:

// Worker 脚本
workerPort.onmessage = (e: MessageEvents) => {
  let images = e.data.images;
  let thumbnails: ArrayBuffer[] = [];
  for (let i = 0; i < images.length; i++) {
    let thumb = compressImage(images[i]);
    thumbnails.push(thumb);
    workerPort.postMessage({
      type: 'progress',
      current: i + 1,
      total: images.length
    });
  }
  workerPort.postMessage({ type: 'complete', thumbnails: thumbnails });
};

场景二:数据加密/解密

对大型文件进行 AES 加密, Worker 分块处理,避免加密期间主线程冻结:

// 分块加密,每块完成后报告进度
for (let offset = 0; offset < fileSize; offset += chunkSize) {
  let chunk = readChunk(fileData, offset, chunkSize);
  let encrypted = encryptChunk(chunk, key);
  encryptedChunks.push(encrypted);
  workerPort.postMessage({
    type: 'progress',
    percent: Math.floor((offset / fileSize) * 100)
  });
}

场景三:实时搜索建议

用户在搜索框中输入时,Worker 在后台对大量历史数据进行模糊匹配,每次匹配完成回传建议列表。通过 terminate() 可以取消正在进行的搜索(用户输入了新的搜索词),实现类似"防抖"的效果:

onSearchChange(keyword: string): void {
  // 取消上一次搜索
  if (this.searchWorker) {
    this.searchWorker.terminate();
  }
  // 启动新搜索
  this.searchWorker = new worker.ThreadWorker('entry/ets/workers/SearchWorker.ets');
  this.searchWorker.postMessage({ keyword: keyword, data: this.historyData });
  this.searchWorker.onmessage = (e) => {
    this.suggestions = e.data.results;
  };
}

场景四:Sensor 数据处理

对于需要持续处理传感器数据(加速度计、陀螺仪)的应用——如运动姿态识别——Worker 可以从传感器获取原始数据,在独立线程中进行降噪、特征提取、模式匹配,然后将识别结果发送给主线程更新 UI:

// Worker 中持续运行的传感器数据分析循环
workerPort.onmessage = (e: MessageEvents) => {
  if (e.data.command === 'start') {
    // 持续采集和分析
    setInterval(() => {
      let rawData = readSensorData();
      let gesture = analyzeGesture(rawData);
      workerPort.postMessage({ type: 'gesture', gesture: gesture });
    }, 50); // 每 50ms 分析一次
  }
};

注意事项与最佳实践

1. 及时清理 Worker

每个 Worker 实例对应一个操作系统线程。不清理的 Worker 会造成线程泄漏——线程本身占用内存和 CPU 时间片,过多未清理的 Worker 会导致系统资源耗尽。在以下三个场景务必调用 terminate()

  • 任务完成后(在 complete 消息回调中)。
  • 组件销毁时(在 aboutToDisappear 中)。
  • 错误发生后(在 onerror 回调中)。

2. 控制消息频率

Worker 之间的消息传递涉及序列化和跨线程通信,有一定的开销。每 1% 的进度发一次消息(总共 100 次)是合理的;每次循环都发消息(30 万次)会导致主线程被消息淹没。

3. 避免传递大量数据

postMessage 会序列化整个消息对象。如果计算结果有 5 万个质数,全部传回主线程会导致严重的序列化开销和内存峰值。Demo 的做法——只传回统计信息(总数、耗时、样本)——是正确的。如果确实需要全部数据,考虑分页传输或写入文件后再通知主线程。

4. Worker 不能直接访问 UI

Worker 脚本运行在独立线程中,@State@Componentbuild() 等 UI 概念在 Worker 中不存在。Worker 唯一能做的 UI 相关操作是通过 postMessage 发送数据,让主线程的代码来更新 UI。

5. 考虑使用 TaskPool

如前面讨论的——如果你只是想做"在后台执行一个函数并获取结果",TaskPool 比 Worker 更简单。不要为了用 Worker 而用 Worker。

总结

本文通过构建一个"并发计算实验室",深入讲解了 HarmonyOS @ohos.worker 多线程 API 的核心用法:

  1. ThreadWorker(path):创建独立线程实例,加载 worker 脚本。每个实例对应一个真实的操作系统线程。
  2. postMessage/onmessage:主线程与 Worker 之间的双向消息通道。建议设计清晰的 type 字段消息协议。
  3. onerror:Worker 异常的统一处理入口。错误发生后应调用 terminate() 清理资源。
  4. terminate():强制终止 Worker 线程并释放系统资源。在 aboutToDisappear 中清理是防止内存泄漏的关键。
  5. 进度报告:Worker 在计算循环中定期发送进度消息,主线程实时更新 UI 进度条,实现流畅的用户体验。

多线程的本质不是让计算变得更快——CPU 的总算力是固定的。多线程的价值在于分离关注点:让计算线程专心计算,让 UI 线程专心渲染,两者互不阻塞。这样即便计算需要 5 秒,用户也能看到一个正在移动的进度条和一个可以点击的"取消"按钮——这就是好的用户体验。

在 HarmonyOS 的并发体系中,@ohos.worker 提供了最大程度的控制力,而 @ohos.taskpool 提供了最大程度的简洁性。根据任务的特点选择合适的工具,你的应用就能在性能和代码复杂度之间找到最佳平衡。


Logo

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

更多推荐