文档管理中心

同步任务开发指导 (TaskPool和Worker)

同步任务通过多个线程之间的协作和同步(如使用锁防止数据竞争),确保任务按特定顺序和规则进行,以保障数据的正确性和程序的正确执行。

当同步任务之间相对独立时,推荐使用TaskPool,例如一系列导入的静态方法或单例实现的方法。如果同步任务之间有关联性,则需要使用Worker。

使用TaskPool处理同步任务

以下场景推荐使用TaskPool。

  • 调度相互独立的任务。

  • 静态方法实现的任务。

  • 单例构造的句柄或者类对象跨线程使用。

说明

由于Actor模型不同线程间内存隔离的特性,非线程安全的单例无法在不同线程间使用。可通过共享模块导出单例解决此问题。

  1. 定义并发函数,实现业务逻辑。

  2. 创建任务Task,通过execute()接口执行该任务。

  3. 对任务返回的结果进行操作。

如下示例中业务使用TaskPool调用相关同步方法的代码,首先定义并发函数taskpoolFunc,需要注意必须使用@Concurrent装饰器装饰该函数;其次定义函数mainFunc,该函数功能为创建任务,执行任务并处理任务返回的结果。

收起
自动换行
深色代码主题
复制
  1. import { taskpool } from '@kit.ArkTS';
  2. // ...
  3. // 步骤1: 定义并发函数,实现业务逻辑
  4. @Concurrent
  5. async function taskpoolFunc(num: number): Promise<number> {
  6. // 根据业务逻辑实现相应的功能
  7. let tmpNum: number = num + 100;
  8. return tmpNum;
  9. }
  10. async function mainFunc(): Promise<void> {
  11. // 步骤2: 创建任务并执行
  12. let task1: taskpool.Task = new taskpool.Task(taskpoolFunc, 1);
  13. let res1: number = await taskpool.execute(task1) as number;
  14. let task2: taskpool.Task = new taskpool.Task(taskpoolFunc, res1);
  15. let res2: number = await taskpool.execute(task2) as number;
  16. // 步骤3: 对任务返回的结果进行操作
  17. console.info(`taskpool: task res1 is: ${res1}`);
  18. console.info(`taskpool: task res2 is: ${res2}`);
  19. }
  20. const MSG_SET = 0;
  21. const MSG_GET = 1;
  22. @Entry
  23. @Component
  24. struct Index {
  25. @State message: string = 'Hello World';
  26. build() {
  27. Row() {
  28. Column() {
  29. Text(this.message)
  30. .fontSize(50)
  31. .fontWeight(FontWeight.Bold)
  32. .onClick(async () => {
  33. await mainFunc();
  34. // ...
  35. })
  36. }
  37. .width('100%')
  38. }
  39. .height('100%')
  40. }
  41. }

使用Worker处理关联的同步任务

当一系列同步任务需要使用同一个句柄调度,或者需要依赖某个类对象调度,且无法在不同任务池之间共享时,需要使用Worker。

  1. 在UI主线程中创建Worker对象并接收Worker线程发送的消息。DevEco Studio支持一键生成Worker。在{moduleName}目录下任意位置,点击鼠标右键 > New > Worker,即可生成Worker的模板文件及配置信息。

    收起
    自动换行
    深色代码主题
    复制
    1. import { MessageEvents, worker } from '@kit.ArkTS';
    2. // ...
    3. @Entry
    4. @Component
    5. struct Index {
    6. @State message: string = 'Hello World';
    7. build() {
    8. Row() {
    9. Column() {
    10. Text(this.message)
    11. .fontSize(50)
    12. .fontWeight(FontWeight.Bold)
    13. .onClick(async () => {
    14. // ...
    15. let w: worker.ThreadWorker = new worker.ThreadWorker('entry/ets/workers/MyWorker2.ets');
    16. w.onmessage = (e: MessageEvents): void => {
    17. // 接收Worker子线程的结果
    18. console.info(`main thread onmessage, ${e.data.message}`);
    19. // 销毁Worker
    20. if (e.data.isTerminate) {
    21. w.terminate();
    22. }
    23. }
    24. // 向Worker子线程发送Set消息
    25. w.postMessage({'type': MSG_SET, 'data': 10});
    26. // 向Worker子线程发送Get消息
    27. w.postMessage({'type': MSG_GET});
    28. })
    29. }
    30. .width('100%')
    31. }
    32. .height('100%')
    33. }
    34. }
  2. 在Worker线程中绑定Worker对象,同时处理同步任务逻辑。

    收起
    自动换行
    深色代码主题
    复制
    1. export default class Handle {
    2. id: number = 0;
    3. syncGet(): number {
    4. return this.id;
    5. }
    6. syncSet(num: number): boolean {
    7. this.id = num;
    8. return true;
    9. }
    10. }
    收起
    自动换行
    深色代码主题
    复制
    1. import { worker, ThreadWorkerGlobalScope, MessageEvents } from '@kit.ArkTS';
    2. // 导入句柄类型
    3. import Handle from './handle';
    4. const MSG_SET = 0;
    5. const MSG_GET = 1;
    6. let workerPort : ThreadWorkerGlobalScope = worker.workerPort;
    7. // 无法传输的句柄,所有操作依赖此句柄
    8. let handler: Handle = new Handle();
    9. // Worker线程的onmessage逻辑
    10. workerPort.onmessage = (e : MessageEvents): void => {
    11. switch (e.data.type as number) {
    12. case MSG_SET:
    13. let result: boolean = handler.syncSet(e.data.data);
    14. console.info('worker: result is ' + result);
    15. workerPort.postMessage({'message': 'the result of syncSet() is ' + result, 'isTerminate': false});
    16. break;
    17. case MSG_GET:
    18. let num: number = handler.syncGet();
    19. console.info('worker: num is ' + num);
    20. workerPort.postMessage({'message': 'the result of syncGet() is ' + num, 'isTerminate': true});
    21. break;
    22. default:
    23. workerPort.postMessage({ 'message': 'send message is invalid', 'isTerminate': false });
    24. break;
    25. }
    26. }