多线程并发TaskPool 和 Worker

异步并发能解决 I/O 等待的问题,但解决不了 CPU 忙的问题。一个图片编解码任务,CPU 在那吭哧吭哧算,你用 async/await 把它挂起,CPU 还是在算那个任务,主线程照样卡。这种要靠多线程——开一个子线程,把计算扔过去,主线程继续响应 UI。ArkTS 提供了两套多线程方案:TaskPool 和 Worker。这篇先把多线程并发的概念讲清楚,再把这两套方案的定位和选择思路说透。

为什么需要多线程并发

先把异步并发和多线程并发的区别讲清楚,这是理解后续内容的基础。

异步并发是指异步代码在执行到一定程度后暂停,并在未来某个时间点继续执行,同一时间只有一段代码执行。ArkTS 通过 Promise 和 async/await 提供异步并发能力,适用于单次 I/O 任务。

多线程并发允许同时执行多段代码。UI 主线程继续响应用户操作和更新 UI,后台线程执行耗时操作,避免应用卡顿。ArkTS 通过 TaskPool 和 Worker 提供多线程并发能力,适用于耗时任务等并发场景。

I/O 等待

CPU 计算

并发需求

任务类型

异步并发
Promise / async-await

多线程并发
TaskPool / Worker

单线程调度
等待时不阻塞

多线程并行
真正同时执行

关键区别在"同一时间只有一段代码执行" vs “同时执行多段代码”。异步并发是单线程下的调度技巧,CPU 还是一个;多线程并发是真的有多个 CPU 核在同时干活。

在多核设备上,任务可以在不同 CPU 上并行执行。对于单核设备,尽管多个任务不会同时执行,但 CPU 会在某个任务休眠或进行 I/O 操作时切换任务,调度其他任务,提高 CPU 的资源利用率。

多线程并发的典型场景

多线程并发适用于多种业务场景,常见的分三类:

长时间执行的任务:业务逻辑包含大量计算或频繁的 I/O 读写,例如图片和视频的编解码、文件的压缩与解压缩、数据库操作等场景。这些任务跑起来要几百毫秒甚至几秒,放主线程 UI 必卡。

长时间保持运行的任务:业务逻辑包括监听和定期采集数据,例如定期采集传感器数据的场景。这种任务要一直跑着,不能放主线程占着。

跟随主线程生命周期的任务:业务逻辑跟随主线程的生命周期,或与主线程绑定的任务,例如游戏中的业务场景。这种任务跟主线程同生共死,但又要独立执行。

Actor 并发模型

讲 TaskPool 和 Worker 之前,得先讲 Actor 并发模型,因为这两个都是基于 Actor 模型实现的。

并发模型用于实现不同应用场景中的并发任务。常见的有基于内存共享的模型和基于消息通信的模型。Actor 并发模型是基于消息通信的典型并发模型,开发者无需处理锁带来的复杂问题,且具备高并发度,因此应用广泛。

内存共享模型 vs Actor 模型

内存共享并发模型:多线程同时执行任务,这些线程依赖同一内存资源并且都有权限访问,线程访问内存前需要抢占并锁定内存的使用权,没有抢占到内存的线程需要等待其他线程释放使用权再执行。

Actor 并发模型:每一个线程都是一个独立 Actor,每个 Actor 有自己独立的内存,Actor 之间通过消息传递机制触发对方 Actor 的行为,不同 Actor 之间不能直接访问对方的内存空间。

Actor 模型

消息

消息

消息

Actor1
独立内存

Actor2
独立内存

Actor3
独立内存

内存共享模型

需要加锁

线程1

共享内存

线程2

线程3

两者的核心区别在内存隔离方式。内存共享模型里所有线程能访问同一块内存,访问时要加锁防止冲突;Actor 模型里每个线程有独立内存,线程间不共享内存,通过消息传递通信。

Actor 模型的好处:不同线程间的内存是隔离的,不会发生线程竞争同一内存资源的情况,无需处理内存上锁问题,从而提高开发效率。坏处是线程间通信要走消息传递,有序列化开销。

用 Actor 模型解决生产者消费者问题

经典的生产者消费者问题,用两种模型对比看。内存共享模型下,生产者和消费者共享一个队列,访问时要加锁:

// 内存共享模型伪代码(仅示意,ArkTS 不这么写)
class BufferQueue {
  public queue: Queue = new Queue();
  public mutex: Mutex = new Mutex();

  add(value: number) {
    if (this.mutex.lock()) {      // 抢锁
      this.queue.push(value);
      this.mutex.unlock();        // 释放锁
    }
  }

  take(): number {
    let res: number = 0;
    if (this.mutex.lock()) {      // 抢锁
      let num = this.queue.pop();
      this.mutex.unlock();        // 释放锁
      res = num;
    }
    return res;
  }
}

// 全局共享内存
let gBufferQueue = new BufferQueue();
// 生产者和消费者都访问 gBufferQueue,靠锁协调

Actor 模型下,生产者和消费者各有自己的内存,通过消息传递结果:

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

// 生产者:在子线程跑
@Concurrent
async function produce(): Promise<number> {
  console.info('producing...');
  return Math.random();  // 生产结果,通过消息传回主线程
}

// 消费者:在主线程跑
class Consumer {
  public consume(value: number) {
    console.info(`consuming value: ${value}`);
  }
}

@Entry
@Component
struct ActorModel {
  @State message: string = 'Hello World';

  build() {
    Column() {
      Button('Actor start')
        .onClick(() => {
          let produceTask: taskpool.Task = new taskpool.Task(produce);
          let consumer: Consumer = new Consumer();
          for (let index = 0; index < 10; index++) {
            // 执行生产任务,结果通过消息传回
            taskpool.execute(produceTask).then((res: Object) => {
              consumer.consume(res as number);  // 消费
            }).catch((e: Error) => {
              console.error(`produceTask failed: ${e.message}`);
            });
          }
        })
    }
  }
}

Actor 模型下没有锁,生产者在自己线程里生产,结果通过 taskpool.execute 的 Promise 返回传给主线程,主线程消费。逻辑清晰,不用操心锁的问题。

TaskPool 的定位

TaskPool 是 ArkTS 提供的轻量级并发方案。定位是:轻量级、自动管理、适合短任务。

核心特点:

  • 任务池机制:底层维护一个线程池,任务来了分配线程执行,不用每次新建线程
  • 自动管理线程:线程的创建、复用、销毁由系统管理,开发者不用管
  • 适合短任务:单次执行时间不长(毫秒到秒级)的任务
  • API 简单:创建 Task,调 execute,拿 Promise 结果
import { taskpool } from '@kit.ArkTS';

// 用 @Concurrent 标记可以在子线程执行的函数
@Concurrent
function imageDecode(data: ArrayBuffer): ArrayBuffer {
  // 图片解码,CPU 密集
  return decodeImpl(data);
}

// 主线程调用
async function decodeImage(data: ArrayBuffer): Promise<ArrayBuffer> {
  let task = new taskpool.Task(imageDecode, data);
  let result = await taskpool.execute(task) as ArrayBuffer;
  return result;
}

TaskPool 适合的场景:图片编解码、数据加解密、JSON 解析、数学计算这类 CPU 密集但单次不太久的任务。任务来了扔给池子,池子分配线程跑,跑完结果回来。

Worker 的定位

Worker 是 ArkTS 提供的重量级并发方案。定位是:重量级、手动管理、适合长任务。

核心特点:

  • 独立线程:每个 Worker 是一个独立的线程,不是池化复用的
  • 手动管理生命周期:开发者自己创建、销毁 Worker
  • 适合长任务:需要长时间运行或与主线程频繁通信的任务
  • 双向通信:主线程和 Worker 之间可以双向发消息
// worker.ts - Worker 线程代码
import { worker } from '@kit.ArkTS';

const workerPort: worker.ThreadWorkerGlobalScope = worker.workerPort;

// 接收主线程消息
workerPort.onmessage = (e: MessageEvents) => {
  let data = e.data;
  // 处理数据
  let result = process(data);
  // 把结果发回主线程
  workerPort.postMessage(result);
};
// 主线程代码
import { worker } from '@kit.ArkTS';

let myWorker = new worker.Worker('entry/ets/workers/worker.ts');

// 接收 Worker 消息
myWorker.onmessage = (e: MessageEvents) => {
  console.info(`收到 Worker 结果: ${e.data}`);
};

// 发消息给 Worker
myWorker.postMessage({ type: 'process', data: someData });

// 用完销毁
myWorker.terminate();

Worker 适合的场景:后台长期运行的数据采集、与主线程频繁交互的复杂计算、需要独立维持状态的长时间任务。

TaskPool vs Worker 对比

把两者的特点放一起对比:

维度TaskPoolWorker
定位轻量级短任务重量级长任务
线程管理自动,池化复用手动,独立线程
生命周期任务执行完即结束开发者控制创建销毁
通信方式单次请求-响应(Promise)双向消息传递
API 复杂度简单(Task + execute)中等(Worker + onmessage + postMessage)
资源开销低(线程复用)较高(独立线程)
并发度系统调度,可多任务并行一个 Worker 一个线程
状态保持无状态,每次独立有状态,可跨消息保持
适合任务时长毫秒到秒级秒级以上,长期运行
典型场景图片解码、数据加密、JSON 解析后台采集、长连接、复杂计算

短任务、无状态

长任务、有状态

需要频繁双向通信

一次性请求-响应

多线程需求

任务特征

TaskPool

Worker

创建 Task

execute 提交

await 等结果

创建 Worker

postMessage 通信

onmessage 处理

terminate 销毁

两者都基于 Actor 模型

TaskPool 和 Worker 都基于 Actor 并发模型实现。这意味着:

  • 线程间内存隔离,不共享内存
  • 通信通过消息传递(序列化)
  • 不用处理锁的问题

但两者在 Actor 模型的实现方式上有差异。TaskPool 的 Actor 是池化的,任务来了分配一个 Actor 执行,执行完归还;Worker 的 Actor 是独占的,一个 Worker 对应一个长期存在的 Actor。

选择思路

怎么在 TaskPool 和 Worker 之间选?按这几个问题走一遍:

任务执行多久? 毫秒到秒级选 TaskPool,秒级以上或长期运行选 Worker。TaskPool 的线程复用优势在短任务上明显,长任务复用价值不大。

需要跟主线程频繁交互吗? 不需要交互,扔过去算完拿结果回来,选 TaskPool。需要多次往返通信,选 Worker。TaskPool 是单次请求-响应,Worker 支持双向消息。

任务需要保持状态吗? 无状态任务选 TaskPool,每次执行独立。有状态任务(比如要维护一个数据结构跨多次操作)选 Worker,Worker 可以跨消息保持状态。

任务数量多吗? 大量同类短任务选 TaskPool,池化复用省开销。少量长任务选 Worker,每个 Worker 独立管理。

资源敏感吗? 内存紧张选 TaskPool,线程池复用省资源。资源充足且任务确实需要长线程,选 Worker。

来看几个具体场景的选择:

// 场景1:列表里每张图片要解码 → TaskPool
// 短任务、无状态、数量多、不需要交互
@Concurrent
function decodeImage(data: ArrayBuffer): ArrayBuffer {
  return imageDecodeImpl(data);
}

async function loadImages(images: ArrayBuffer[]): Promise<ArrayBuffer[]> {
  let tasks = images.map(img => {
    let task = new taskpool.Task(decodeImage, img);
    return taskpool.execute(task);
  });
  return await Promise.all(tasks) as ArrayBuffer[];
}

// 场景2:后台持续采集传感器数据 → Worker
// 长任务、有状态、需要持续通信
// worker.ts
// workerPort.onmessage = (e) => {
//   if (e.data.type === 'start') {
//     startSensorLoop();  // 持续采集,定期 postMessage 回主线程
//   }
// };

// 场景3:单次大文件压缩 → TaskPool
// 短任务(相对而言)、无状态、一次性
@Concurrent
function compressFile(data: ArrayBuffer): ArrayBuffer {
  return compressImpl(data);
}

// 场景4:游戏里的 AI 逻辑 → Worker
// 长任务、有状态、需要频繁跟主线程交互
// AI Worker 持续运行,主线程发玩家操作,Worker 算 AI 决策

并发注意事项

不管用 TaskPool 还是 Worker,有几个共通的坑要注意:

避免在并发线程中操作 UI。 UI 操作必须在主线程中执行。并发线程中操作 UI 可能导致界面异常或崩溃。子线程算完数据,把数据传回主线程,主线程再刷新 UI。

数据传递需支持序列化。 并发任务间传递数据时,对象必须是可序列化的(如基本类型、普通对象等),或者可共享的(sendable 对象),不可传递函数、循环引用、特殊对象(如 Promise、Error)等。已完成(fulfilled 或 rejected)状态的 Promise 可以被传递,因为其结果是可序列化的。

合理控制并发粒度。 频繁创建和销毁并发任务(如 Worker、Task)会带来额外性能开销,建议复用或使用任务池机制。TaskPool 天然复用,Worker 要自己管好别频繁创建销毁。

注意内存泄漏风险。 避免在并发任务中持有外部对象的强引用,防止内存泄漏。Worker 用完要 terminate,别让没用的 Worker 一直占着线程。

并发任务应具备独立性。 并发任务应尽量不依赖外部状态,减少竞态条件和同步开销。竞态条件是指多个线程或任务同时访问并修改共享数据,执行结果依赖于任务调度的顺序,可能导致数据不一致或不可预期的行为。Actor 模型下内存隔离,但通过消息传递的状态还是要设计好。

一个完整的选型案例

场景:一个相册应用,要加载 20 张图片的缩略图,同时后台还要持续分析图片的人脸数据。

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

// 缩略图加载:短任务、无状态、数量多 → TaskPool
@Concurrent
function generateThumbnail(imageData: ArrayBuffer): ArrayBuffer {
  // 图片缩放解码,毫秒级
  return thumbnailImpl(imageData);
}

// 人脸分析:长任务、有状态、持续运行 → Worker
// face-worker.ts
// const workerPort = worker.workerPort;
// let faceDetector = new FaceDetector();  // 跨消息保持状态
// workerPort.onmessage = (e) => {
//   if (e.data.type === 'detect') {
//     let result = faceDetector.detect(e.data.image);
//     workerPort.postMessage({ type: 'result', faces: result });
//   }
// };

@Entry
@Component
struct AlbumPage {
  @State thumbnails: ArrayBuffer[] = [];
  @State faceCount: number = 0;
  private faceWorker: worker.Worker | null = null;

  aboutToAppear() {
    // 启动人脸分析 Worker
    this.faceWorker = new worker.Worker('entry/ets/workers/face-worker.ts');
    this.faceWorker.onmessage = (e: MessageEvents) => {
      if (e.data.type === 'result') {
        this.faceCount = e.data.faces.length;
      }
    };
  }

  aboutToDisappear() {
    // 别忘了销毁
    this.faceWorker?.terminate();
  }

  async loadThumbnails(images: ArrayBuffer[]) {
    // 20 张缩略图并行加载,用 TaskPool
    let tasks = images.map(img => {
      let task = new taskpool.Task(generateThumbnail, img);
      return taskpool.execute(task);
    });
    this.thumbnails = await Promise.all(tasks) as ArrayBuffer[];
  }

  detectFaces(image: ArrayBuffer) {
    // 发给 Worker 持续分析
    this.faceWorker?.postMessage({ type: 'detect', image: image });
  }

  build() {
    Column() {
      Text(`检测到 ${this.faceCount} 张人脸`)
      // 缩略图列表...
    }
  }
}

这个案例里两种方案各司其职:缩略图加载用 TaskPool,因为是一次性批量短任务;人脸分析用 Worker,因为要持续运行且保持检测器状态。混着用才是常态,不是非此即彼。

总结一下下

先想清楚任务特征再选方案。 别上来就 Worker,大部分场景 TaskPool 就够。Worker 的管理成本高,能不用就不用。判断标准:要不要长期运行、要不要双向通信、要不要保持状态,三个都否就用 TaskPool。

TaskPool 任务要轻。 TaskPool 的设计假设是短任务,单个任务执行太久会占着池子里的线程,影响其他任务。如果任务要跑几秒以上,考虑 Worker。

Worker 创建销毁有成本。 别频繁创建销毁 Worker,每个 Worker 是独立线程,创建有开销。需要反复用的场景,创建一次一直用,用完再销毁。

消息传递的数据要序列化。 跨线程传的数据会被序列化,大对象传递有开销。能传引用的用 Sendable 对象,不能的尽量传小数据。别把整个图片 ArrayBuffer 频繁传来传去。

子线程异常不会自动到主线程。 TaskPool 的异常通过 Promise 的 reject 传回,要 catch;Worker 的异常要监听 onerror。别假设子线程不会出错,异常处理要跟上。

UI 刷新一定回主线程。 子线程算完数据,把数据传回主线程,由主线程更新 @State 变量触发刷新。直接在子线程改 UI 状态会出问题。

Logo

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

更多推荐