鸿蒙多线程开发实践与避坑指南

前面几篇讲了 TaskPool、Worker、序列化、Sendable 的原理和用法。这篇把这些东西串起来,回到真实业务场景:图片处理、数据计算、文件 IO、网络请求这些常见耗时任务,到底怎么用多线程写,怎么避坑。每个场景给完整代码示例,再总结调试方法和常见错误。

常见多线程场景

ArkTS 多线程并发主要解决三类问题:

耗时任务

大量计算
图片编解码 / 视频编码 / 数据分析

频繁 I/O
文件读写 / 数据库操作 / 网络请求

长时运行
传感器采集 / 算法训练 / 常驻任务

TaskPool
3 分钟内独立任务

Worker
超 3 分钟或需保存状态

  • 大量计算:图片编解码、视频编码、数据分析、加密解密。CPU 密集型,会阻塞线程。
  • 频繁 I/O:文件读写、数据库操作、网络请求。I/O 密集型,异步等待时间长。
  • 长时运行:传感器数据采集、算法训练、常驻后台任务。需要持续占用线程。

3 分钟以内的独立任务用 TaskPool,超过 3 分钟或需要保存线程上下文的用 Worker。

TaskPool 实战:批量图片处理

场景:相册里批量给图片做滤镜处理,每张图片独立处理,结果回主线程显示。

分析:每张图片处理几秒,独立任务,需要优先级(用户当前看的图先处理),可能要取消(用户滑走了就别处理了)。完美匹配 TaskPool。

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

// 滤镜处理函数,必须在单独文件或用 @Concurrent 标注
@Concurrent
function applyFilter(buffer: ArrayBuffer, filterType: number): ArrayBuffer {
  const pixels = new Uint8Array(buffer);
  // 简化:根据 filterType 做不同处理
  for (let i = 0; i < pixels.length; i += 4) {
    if (filterType === 1) {
      // 灰度
      const gray = pixels[i] * 0.299 + pixels[i + 1] * 0.587 + pixels[i + 2] * 0.114;
      pixels[i] = gray;
      pixels[i + 1] = gray;
      pixels[i + 2] = gray;
    } else if (filterType === 2) {
      // 反色
      pixels[i] = 255 - pixels[i];
      pixels[i + 1] = 255 - pixels[i + 1];
      pixels[i + 2] = 255 - pixels[i + 2];
    }
  }
  return buffer;
}

// 批量处理多张图片
async function batchProcessImages(images: { buffer: ArrayBuffer; filterType: number }[]): Promise<ArrayBuffer[]> {
  const group = new taskpool.TaskGroup();
  for (const img of images) {
    // 用 setCloneList 避免转移所有权,因为 buffer 可能要复用
    const task = new taskpool.Task(applyFilter, img.buffer, img.filterType);
    task.setCloneList([img.buffer]);
    group.addTask(task);
  }
  try {
    // HIGH 优先级,图片处理影响用户体验
    const results = await taskpool.execute(group, taskpool.Priority.HIGH) as ArrayBuffer[];
    return results;
  } catch (err) {
    const error = err as BusinessError;
    console.error(`batchProcess failed: ${error.code}, ${error.message}`);
    return [];
  }
}

// 取消场景:用户滑走了,取消未处理的任务
let currentTask: taskpool.Task | null = null;

async function processSingleWithCancel(buffer: ArrayBuffer, filterType: number): Promise<ArrayBuffer | null> {
  // 取消上一个未完成的任务
  if (currentTask) {
    try {
      currentTask.cancel();
    } catch (e) {
      // 任务可能已经完成,忽略
    }
  }
  currentTask = new taskpool.Task(applyFilter, buffer, filterType);
  try {
    const result = await taskpool.execute(currentTask, taskpool.Priority.HIGH) as ArrayBuffer;
    currentTask = null;
    return result;
  } catch (err) {
    currentTask = null;
    if ((err as BusinessError).message.includes('canceled')) {
      return null; // 取消了,返回 null
    }
    throw err;
  }
}

几个要点:

  • 用 TaskGroup 批量提交,结果数组顺序和提交顺序一致。
  • setCloneList 让 ArrayBuffer 走拷贝,避免转移后失效。如果传完不再用,去掉这行用转移更快。
  • HIGH 优先级,图片处理直接影响 UI 显示。
  • 取消用 task.cancel(),取消后 await 会抛异常,message 含 canceled。

Worker 实战:长时间数据计算

场景:用房价数据训练预测模型,训练要 10 分钟以上,训练完用模型预测。模型要保存在 Worker 线程里,预测依赖训练结果。

分析:超过 3 分钟,需要保存模型状态(句柄),强关联同步任务。完美匹配 Worker。

主线程:

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

class PricePredictor {
  private workerInstance: worker.ThreadWorker;

  constructor() {
    this.workerInstance = new worker.ThreadWorker('entry/ets/workers/PredictWorker.ets');

    this.workerInstance.onAllErrors = (err: ErrorEvent) => {
      console.error(`Worker error: ${err.message}`);
    };

    this.workerInstance.onexit = (code: number) => {
      console.info(`Worker exit, code: ${code}`);
    };
  }

  // 训练模型
  train(): Promise<string> {
    return new Promise((resolve) => {
      this.workerInstance.onmessage = (e: MessageEvents) => {
        if (e.data.type === 'train') {
          resolve(e.data.value);
        }
      };
      this.workerInstance.postMessage({ type: 0 });
    });
  }

  // 预测
  predict(area: number, room: number): Promise<number> {
    return new Promise((resolve) => {
      this.workerInstance.onmessage = (e: MessageEvents) => {
        if (e.data.type === 'predict') {
          resolve(e.data.value);
        }
      };
      this.workerInstance.postMessage({ type: 1, area, room });
    });
  }

  // 销毁
  destroy() {
    this.workerInstance.terminate();
  }
}

// 使用
async function runPrediction() {
  const predictor = new PricePredictor();
  try {
    const trainResult = await predictor.train();
    console.info(`训练完成: ${trainResult}`);

    const price = await predictor.predict(80, 4);
    console.info(`预测房价: ${price} 元/平米`);
  } catch (err) {
    console.error(`预测失败: ${(err as BusinessError).message}`);
  } finally {
    predictor.destroy(); // 用完销毁
  }
}

PredictWorker.ets:

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

const workerPort: ThreadWorkerGlobalScope = worker.workerPort;

class PriceModel {
  areaCoefficient: number = 0;
  roomCoefficient: number = 0;
  basePrice: number = 0;
}

// 模型保存在 Worker 线程上下文里
const model: PriceModel = new PriceModel();

function predict(area: number, room: number): number {
  return (model.areaCoefficient * area + model.roomCoefficient * room) * model.basePrice;
}

function optimize(): void {
  // 简化的训练过程,实际可能要跑很久
  model.areaCoefficient = 3;
  model.roomCoefficient = 500;
  model.basePrice = 10;
}

workerPort.onmessage = (e: MessageEvents): void => {
  switch (e.data.type as number) {
    case 0: // 训练
      optimize();
      workerPort.postMessage({ type: 'train', value: 'train success' });
      break;
    case 1: // 预测
      const output = predict(e.data.area as number, e.data.room as number);
      workerPort.postMessage({ type: 'predict', value: output });
      break;
    default:
      workerPort.postMessage({ type: 'message', value: 'invalid type' });
      break;
  }
};

要点:

  • Worker 实例复用,不要每次训练都 new 一个。
  • 模型保存在 Worker 线程的全局上下文里,多次消息传递共享这个状态。这是 Worker 的核心优势。
  • 用完 terminate() 销毁,释放内存。
  • 主线程用 Promise 包装消息通信,调用方代码更干净。

大数据传输优化

场景:子线程从数据库读 5MB 数据返回主线程渲染。

错误做法:直接 return,走序列化。5MB 序列化反序列化要几十毫秒,还可能撞 16MB 限制。

// 不推荐:序列化传输 5MB 数据
@Concurrent
function loadDataBad(): ArrayBuffer {
  // 从数据库读 5MB
  return bigBuffer; // 序列化开销大
}

正确做法 1:用 Sendable 共享。

import { collections } from '@arkts.collections';

@Sendable
class SharedData {
  data: collections.Array<number>;
  constructor(data: collections.Array<number>) {
    this.data = data;
  }
}

@Concurrent
async function loadDataSendable(): Promise<SharedData> {
  // 从数据库读数据,构造 Sendable 对象
  const arr = new collections.Array<number>();
  for (let i = 0; i < 5000000; i++) {
    arr.push(Math.random());
  }
  return new SharedData(arr); // 引用传递,不拷贝
}

async function useSendable() {
  const task = new taskpool.Task(loadDataSendable);
  const shared = await taskpool.execute(task) as SharedData;
  // 直接用 shared.data,和子线程引用同一个对象
  console.info(`data size: ${shared.data.length}`);
}

正确做法 2:用 ArrayBuffer 转移所有权(一次性传输)。

@Concurrent
function loadDataTransfer(): ArrayBuffer {
  // 从数据库读到 ArrayBuffer
  return bigBuffer;
}

async function useTransfer() {
  const task = new taskpool.Task(loadDataTransfer);
  const buffer = await taskpool.execute(task) as ArrayBuffer;
  // buffer 所有权已转移到主线程,子线程那边失效
  // 注意:如果子线程还要用这个 buffer,不能用转移
}

选择:

  • 子线程和主线程都要访问 → Sendable
  • 只传一次,传完一边不用 → ArrayBuffer 转移
  • 数据量小(< 100KB)→ 直接序列化,简单

内存管理

多线程内存管理几个关键点:

及时销毁 Worker。 Worker 空闲也占内存,不用了立刻 terminate()。Worker 数量上限 64 个,加上 napi runtime 不超过 80,超了直接报错。所有 Worker + 主线程累积内存超 1.5GB 或设备内存 60% 会 OOM 崩溃。

TaskPool 缩容要满足条件。 TaskPool 30 秒检测一次,空闲线程要满足 5 个条件才释放(空闲 30 秒、无 LongTask、无未释放句柄、非调试、无子 Worker)。如果任务里创建了 Timer 没释放,对应线程不会缩容。

避免内存泄漏。 并发任务里别持有外部对象的强引用。Worker 的 onmessage 回调里如果捕获了主线程的大对象,会一直引用着不释放。

// 错误:onmessage 闭包持有大对象
let bigData = new ArrayBuffer(100 * 1024 * 1024); // 100MB
workerInstance.onmessage = (e: MessageEvents) => {
  console.info(bigData.byteLength); // bigData 被引用,无法释放
};

// 正确:用完置 null,或者用弱引用
workerInstance.onmessage = (e: MessageEvents) => {
  console.info(e.data); // 只用消息里的数据
};
bigData = null; // 用完释放

Sendable 对象生命周期。 Sendable 对象在共享堆上,所有线程都不再引用时才回收。注意别让某个线程长期持有不用的 Sendable 对象引用。

调试多线程代码

多线程调试比单线程麻烦,几个方法:

用 hilog 打日志。 hilog 是线程安全的,多线程都能用。给日志加线程标识便于区分:

import { hilog } from '@kit.BasicServicesKit';

const TAG = 'WorkerTag';

@Concurrent
function someTask() {
  hilog.info(0x0000, TAG, 'task start'); // 工作线程打日志
  // ...
  hilog.info(0x0000, TAG, 'task end');
}

捕获 Worker 全局异常。 主线程注册 onAllErrors 能捕获 Worker 线程里 onmessage、timer 回调、文件执行等流程的全局异常:

workerInstance.onAllErrors = (err: ErrorEvent) => {
  console.error(`Worker error: ${err.message}, stack: ${err.stack}`);
};

监控任务耗时。 TaskPool 的 Task 有 cpuDuration 和 ioDuration 属性,可以看 CPU 执行时长和异步 I/O 等待时长,判断是否接近 3 分钟限制。

const task = new taskpool.Task(someFunc, ...);
await taskpool.execute(task);
console.info(`CPU: ${task.cpuDuration}ms, IO: ${task.ioDuration}ms`);

DevEco Studio Profiler。 用 Profiler 看线程状态、内存占用、CPU 使用率,能直观看到 Worker 创建销毁、TaskPool 扩缩容过程。

常见错误和解决方案

错误 1:@Concurrent 函数调用同文件函数

function helper() { }

@Concurrent
function task() {
  helper(); // 报错:Only imported variables and local variables can be used
}

解决:把 helper 移到单独文件 import 进来,或者把逻辑直接写在 task 里。

错误 2:任务超过 3 分钟被阻塞

@Concurrent
function longTask() {
  while (true) {
    // 死循环,超过 3 分钟,线程被占满
  }
}

解决:长任务用 LongTask,或者改用 Worker。如果是同步 I/O 阻塞,改成异步 I/O(不计入 3 分钟限制)。

错误 3:序列化数据超过 16MB

@Concurrent
function process(bigData: ArrayBuffer): ArrayBuffer {
  return bigData;
}

const huge = new ArrayBuffer(20 * 1024 * 1024); // 20MB
const task = new taskpool.Task(process, huge);
await taskpool.execute(task); // 报错:序列化数据超限

解决:用 ArrayBuffer 转移所有权(setTransferList),或用 Sendable 共享。

错误 4:Worker 数量超过 64

const workers: worker.ThreadWorker[] = [];
for (let i = 0; i < 100; i++) {
  workers.push(new worker.ThreadWorker('...')); // 报错:exceeds the maximum
}

解决:用 Worker 池复用,或者改用 TaskPool(自动管理线程池)。

错误 5:Worker 销毁后调用接口

workerInstance.terminate();
workerInstance.postMessage('hello'); // 报错:Worker 已销毁

解决:在 onexit 回调里处理销毁后的逻辑,销毁前确保所有 postMessage 已发出。

错误 6:在 Worker 里操作 UI

// worker.ets
workerPort.onmessage = (e: MessageEvents) => {
  // 试图更新 UI,崩溃
  someComponent.message = 'updated';
};

解决:Worker 算完结果 postMessage 给主线程,主线程在 onmessage 里更新 UI。

错误 7:Sendable 对象数据竞争

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

const counter = new Counter();

@Concurrent
function inc(counter: Counter) {
  counter.count++; // 多个线程同时调,数据竞争
}

解决:用 AsyncLock 加锁,或者把对象冻结成只读。

错误 8:Promise 跨线程传递失败

@Concurrent
function badTask(): Promise<number> {
  return new Promise((resolve) => {
    setTimeout(() => resolve(1), 1000);
  }); // 返回 pending 状态的 Promise,TaskPool 报错
}

解决:用 async/await 等待 Promise 完成:

@Concurrent
async function goodTask(): Promise<number> {
  return await new Promise((resolve) => {
    setTimeout(() => resolve(1), 1000);
  });
}

总结一下下

先量后优。 多线程不是银弹,加线程有创建开销、序列化开销、内存开销。先用单线程跑,确认是性能瓶颈再加多线程。Profiler 数据驱动,别凭感觉。

任务粒度别太细。 每个任务有创建、序列化、调度开销。任务太细(比如每个任务处理一个像素),开销可能比计算还大。把任务合并到合理粒度,比如每个任务处理一行或一块。

复用 Worker。 Worker 创建开销大,别频繁创建销毁。做 Worker 池,按需借还。TaskPool 自动管,没这个问题。

优先级要会用。 影响用户体验的任务用 HIGH,后台同步用 IDLE。中载下优先级效果不明显,重载下高优先级才能抢到资源。

监控内存。 长跑应用定期看内存占用。Worker 数量、TaskPool 线程数、Sendable 对象引用,都要关注。内存涨上去不下来,多半是泄漏。

测试要覆盖多线程场景。 单线程测试通过不代表多线程没问题。数据竞争、死锁、时序问题只在多线程下暴露。写测试时模拟并发场景,多跑几轮。

别在并发线程里用 AppStorage。 AppStorage 只能在主线程用,TaskPool 和 Worker 都不支持。要共享状态用 Sendable 或者消息传递。

应用切后台 Worker 暂停。 应用挂起后 Worker 线程暂停。如果任务有时效要求(比如倒计时),注意这个行为,必要时用后台任务保护。

多线程并发写好了性能提升明显,写不好就是一堆坑。把 TaskPool、Worker、序列化、Sendable 的能力和限制都摸清,按场景选对方案,再加上合理的内存管理和调试手段,大部分问题都能搞定。

Logo

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

更多推荐