【HarmonyOS开发小实践】鸿蒙多线程开发实践与避坑指南
鸿蒙多线程开发实践与避坑指南
前面几篇讲了 TaskPool、Worker、序列化、Sendable 的原理和用法。这篇把这些东西串起来,回到真实业务场景:图片处理、数据计算、文件 IO、网络请求这些常见耗时任务,到底怎么用多线程写,怎么避坑。每个场景给完整代码示例,再总结调试方法和常见错误。
常见多线程场景
ArkTS 多线程并发主要解决三类问题:
- 大量计算:图片编解码、视频编码、数据分析、加密解密。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 的能力和限制都摸清,按场景选对方案,再加上合理的内存管理和调试手段,大部分问题都能搞定。
更多推荐



所有评论(0)