异步任务

异步任务(Async Task)用于把耗时代码——发送邮件、生成报表、推送消息等——从请求处理流程中剥离,交由独立的 Task Worker 进程执行,避免阻塞 Worker 响应请求。Viswoole 基于 Swoole Task 机制封装了任务管理器(TaskManager),提供主题注册、异步投递、队列持久化、有效期淘汰与同步等待结果等能力。

核心组件

组件全限定名职责
TaskManagerViswoole\Core\Server\TaskManager任务主题的注册、投递与分发
TaskProxyViswoole\Core\Server\TaskProxy任务代理对象,作为处理器入参
Task 门面Viswoole\Core\Facade\Task静态代理 TaskManager 的核心方法
TaskServiceViswoole\Core\Service\TaskService任务服务提供者,负责装载任务管理器

TaskService 已在框架默认 config/app.phpservices[] 中注册,任务能力开箱即用。业务代码统一通过 Task 门面调用:

php
use Viswoole\Core\Facade\Task;

Task::register('email', SendEmailTask::class);      // 注册任务主题
Task::emit('email.notify', $data);                  // 异步投递,立即返回
$result = Task::emitWait('email.send', $data, 0.5); // 同步等待结果

前置配置

任务能力要求服务配置中 Constant::OPTION_TASK_WORKER_NUM 大于 0(任务进程数量,0 表示不开启)。框架默认 config/server.php 已按 CPU 核数开启任务进程,一般无需修改:

php
// config/server.php
use Swoole\Constant;

return [
  'servers' => [
    'http' => [
      // ...
      'options' => [
        // 任务进程数量,0 表示不开启任务能力
        Constant::OPTION_TASK_WORKER_NUM => swoole_cpu_num(),
      ],
    ],
  ],
];

框架在合并服务配置时会强制开启 Constant::OPTION_TASK_ENABLE_COROUTINE(任务协程化),无需手动配置。

不要自行注册 onTask 事件

框架的任务管理器在实例化时已经通过 ServerEventHook 注册了 task 事件分发器,请勿再通过 ServerEventHook::addEvent('task', ...) 或配置文件注册同名的 task 处理器:

  • Swoole 对同一事件重复注册时会覆盖前一次设置,自行注册会替换框架的任务分发器,导致任务系统整体失效;
  • Swoole 全部服务端事件中仅 onTask 的返回值有语义(将任务结果回传给 Worker 进程),即使多个处理器共存,也只有最后一个处理器的返回值生效,会破坏 emitWait() 的同步等待结果。

任务的业务逻辑一律通过 Task::register() 注册主题实现。详见 生命周期钩子

注册任务

任务按「主题(Topic)」组织:投递时指定主题名,框架将任务分发给该主题注册的处理器。主题名不区分大小写。

注册单个主题(回调方式)

php
use Viswoole\Core\Facade\Task;
use Viswoole\Core\Server\TaskProxy;

Task::register('email.send', function (TaskProxy $task, \Swoole\Server $server): void {
  $data = $task->data; // 投递时传入的业务数据

  // ...执行发送邮件逻辑
  $messageId = sendMail($data['to'], $data['subject'], $data['body']);

  // 回传结果给 Worker 进程(配合 emitWait 使用)
  $task->finish(['ok' => true, 'message_id' => $messageId]);
});

注册类批量注册

register() 的第二个参数传入类名时,框架通过反射扫描该类的全部公开方法,以 「主题前缀.方法名」 的格式批量注册:

php
namespace App\Task;

use Viswoole\Core\Server\TaskProxy;

class SendEmailTask
{
  /**
   * 异步通知:投递后立即返回,无需回传结果
   */
  public function notify(TaskProxy $task): void
  {
    // ...发送通知邮件
  }

  /**
   * 同步等待场景:必须 return 结果,否则调用方会阻塞到超时
   */
  public function send(TaskProxy $task): array
  {
    // ...发送邮件并返回结果
    return ['ok' => true, 'message_id' => '20260912001'];
  }
}
php
// 注册后生成 email.notify、email.send 两个主题
Task::register('email', \App\Task\SendEmailTask::class);

// 投递
Task::emit('email.notify', ['to' => 'user@example.com']);

注册规则:

  • 静态方法注册回调为 ClassName::methodName,实例方法注册为 [ClassName, 'methodName']
  • __ 开头的魔术方法(如 __construct)会被过滤,不会注册为主题;
  • 同一主题重复注册时,后注册的处理器覆盖先注册的。

注册时机

任务主题映射保存在各进程的内存中。请在服务启动阶段注册——例如服务提供者的 register() / boot() 中,或监听 AppInitialized 事件——确保投递任务的 Worker 进程与执行任务的 Task Worker 进程都持有主题映射。参见 服务提供者

异步投递:emit()

emit() 投递任务后立即返回,不等待执行结果,是最常用的投递方式:

参数类型默认值说明
$topicstring已注册的任务主题名称,未注册时抛出 InvalidArgumentException
$datamixed传递给处理器的业务数据
$queuebooltrue是否持久化到队列,启用后服务重启可自动恢复未消费的任务

返回值:投递成功返回任务 ID(int),投递失败返回 false$queue = true 且队列缓存写入失败时抛出 RuntimeException

php
use Viswoole\Core\Facade\Task;

$taskId = Task::emit('email.notify', [
  'to'      => 'user@example.com',
  'subject' => '欢迎注册',
]);

// 明确不持久化的临时任务(服务重启后丢失)
Task::emit('cleanup', ['type' => 'temp'], queue: false);

队列持久化与重启恢复

$queue = true(默认)时,任务在投递前会写入缓存队列(按投递 Worker 进程分标签存储,标签为 $_TASK_QUEUE_ + workerId),随后再投递给 Task Worker。其语义为待消费队列(at-most-once)

  • 任务被 Task Worker 消费时(无论执行成功、失败还是抛异常)即从队列删除,与执行结果无关;
  • 框架在 Worker 进程启动时(workerStart)自动扫描未消费的队列任务并重新投递,因此崩溃恢复覆盖的窗口是**「已投递但尚未被消费」**:Worker 在消费前崩溃(如 kill -9),重启后任务自动恢复执行;
  • 任务执行中途进程崩溃不会重试。如果业务需要「执行失败重试」,应在处理器内部自行实现(如捕获异常后重新 emit)。

队列恢复依赖缓存

队列持久化使用缓存通道存储,缓存数据丢失(如文件缓存被清理、Redis 未持久化)时未消费的任务也会丢失。重要任务建议使用 Redis 等可靠的缓存通道。

任务有效期

验证码短信等任务具有时效性,服务重启后恢复的过期任务再执行只会浪费资源甚至产生副作用。任务有效期用于在消费时淘汰过期任务,在 config/task.php 中配置:

配置键类型默认值说明
storestring|nullnull任务队列持久化使用的缓存商店名称,需在 config/cache.phpstores 中定义;null 表示跟随 cache.default 默认通道
expireint|nullnull全局默认有效期(秒),从投递时刻起算;0null 表示长期有效
topicsarray<string, int|null>[]按主题覆盖有效期(秒),优先级高于 expire,主题名不区分大小写;值为 null 表示该主题豁免(长期有效)
php
// config/task.php
return [
  // 队列读写频繁,默认通道为文件缓存时建议指定 redis 等高性能通道
  'store' => env('task.store'),

  // 全局默认有效期(秒)
  'expire' => env('task.expire'),

  // 按主题覆盖,优先级高于 expire
  'topics' => [
    'email.sendLoginCode' => 300, // 验证码任务 5 分钟有效
    'report.generate'     => 3600,
    'log.write'           => null, // 日志任务永不过期
  ],
];

生效机制:

  • 有效期在 emit() 投递时按主题解析(task.topics 优先于 task.expire),随任务数据一起传递(含持久化到队列的数据);
  • 任务被 Task Worker 消费时检测,已过期的任务不执行处理器,直接向 Worker 进程返回 false 并记录一条任务日志;
  • 仅对 emit() 投递的任务生效;emitWait() / emitsWait() 为同步执行,不受此配置影响;
  • 过期跳过的任务同样遵循队列的 at-most-once 语义,消费时即从队列删除,不会重投。

商店名称必须已定义

store 指向的商店必须已在 cache.stores 中定义,否则首次使用任务队列时会抛出 CacheErrorException(快速失败)。切换商店前残留在旧商店中的队列条目不会被新商店恢复,切换时请确保队列已清空。缓存商店配置见 缓存驱动

同步等待结果:emitWait()

emitWait() 投递任务并阻塞等待执行结果,处理器通过 finish()return 回传的数据类型原样透传

参数类型默认值说明
$topicstring已注册的任务主题名称
$datamixed传递给处理器的业务数据
$timeoutfloat0.5等待超时时间(秒)

返回值:处理器回传的原始结果;主题未注册时抛出 InvalidArgumentException;超时或投递失败返回 false

php
$result = Task::emitWait('email.send', ['to' => 'user@example.com'], 3.0);

if ($result !== false) {
  // $result 为处理器回传的数据,可能是数组、字符串等任意类型
} else {
  // 超时、投递失败,或处理器回传了 false
}

处理器必须显式回传结果

以下为 Swoole taskwait 的实测语义,编写处理器时务必注意:

处理器行为emitWait 结果
finish($data) / return $data(非 false/null)立即返回 $data,类型原样保留
finish(null)立即返回 null(null 作为真实结果送达)
return false / finish(false)立即返回 false与超时的 false 无法区分
无返回值 / return null阻塞满整个超时时间后返回 false(性能陷阱)
任务超时 / 投递失败等待满超时后返回 false

因此:处理器必须显式回传结果,无返回值会导致调用方白白阻塞整个超时时长;表达「业务失败」时应回传结构化数据(如 ['ok' => false])而非裸 false / null

并发等待:emitsWait()

同时投递多个任务并等待全部完成:

参数类型默认值说明
$tasksarray<string, mixed>任务列表,键为主题名称,值为业务数据
$timeoutfloat0.5等待超时时间(秒)
$isCoboolfalsetrue 使用协程并发调度(taskCo),false 使用 taskWaitMulti

返回值:各任务的执行结果列表,超时未完成的任务对应位置为 false

php
$results = Task::emitsWait([
  'email.send'   => ['to' => 'user@example.com'],
  'report.build' => ['user_id' => 100],
], timeout: 3.0, isCo: true);

编写任务处理器

处理器签名

处理器的统一签名为 function (TaskProxy $task, \Swoole\Server $server),第二个参数为 Swoole 服务实例,一般无需使用:

php
function (TaskProxy $task, \Swoole\Server $server): mixed {}

TaskProxy 属性

属性类型默认值说明
datamixednull投递时传入的业务数据
topicstring任务主题名称
queue_idstring|nullnull队列唯一标识,非队列任务为 null

除上述属性外,TaskProxy 还通过 __get 代理访问 Swoole 原生 Task 对象的属性,如 id(任务 ID)、worker_id(所在 Worker 进程 ID)、dispatch_time(投递时间戳)、flags(标志位);访问不存在的属性会抛出异常。

回传结果:return 与 finish()

方式一:处理器 return——返回值经框架原样透传给 Worker 进程(触发 onFinish 事件或 taskwait 返回)。返回 null(含 void 处理器)不会触发结果投递,异步任务(配合 emit())通常直接用这种方式:

php
Task::register('email.query', function (TaskProxy $task): array {
  return ['id' => $task->data['id'], 'status' => 'sent'];
});

方式二:调用 finish(mixed $data): bool——通过 __call 直接透传给底层 Swoole Task 对象,遵循 Swoole 语义可以在任务处理器中多次调用,向 Worker 进程发送多个处理结果(流式任务可逐步回传中间结果):

php
Task::register('report.build', function (TaskProxy $task): void {
  $task->finish(['progress' => 50]); // 中间进度
  // ...继续计算
  $task->finish(['progress' => 100]); // 最终结果
});

注意:队列条目的清理发生在任务被消费时,与 finish() 无关——finish() 只负责结果回传。

异常处理

框架在任务分发时已对 Throwable 做兜底:处理器抛出的异常会被捕获并记录到任务日志通道(含主题、队列 ID、异常原文和任务数据),同时向 Worker 进程返回 falseemitWait 调用方立即收到失败信号而不是干等超时),不会导致 Task Worker 进程崩溃。

业务层面仍建议自行捕获异常并回传结构化结果——裸 false 无法与超时/过期区分:

php
Task::register('email.risky', function (TaskProxy $task): void {
  try {
    $result = doSomethingRisky($task->data);
    $task->finish(['ok' => true, 'data' => $result]);
  } catch (\Throwable $e) {
    $task->finish(['ok' => false, 'error' => $e->getMessage()]);
  }
});

数据序列化

跨进程传输大数据时,可使用 TaskProxy 的静态方法手动序列化/反序列化:

php
$packed = TaskProxy::pack($largeData);  // 序列化为二进制字符串,失败返回 false
$data   = TaskProxy::unpack($packed);   // 反序列化,失败返回 false

完整示例

控制器中注册用户后异步发送欢迎邮件,不阻塞响应:

php
use Viswoole\Core\Facade\Task;
use Viswoole\HttpServer\AutoInject\InjectPost;
use Viswoole\Router\Annotation\{Controller, RouteMapping};

#[Controller(prefix: 'user')]
class UserController
{
  #[RouteMapping(method: 'POST', title: '注册用户')]
  public function register(#[InjectPost] string $email): array
  {
    // ...创建用户

    // 异步发送欢迎邮件(queue 默认 true,服务重启未消费的任务会自动恢复)
    Task::emit('email.notify', [
      'to'      => $email,
      'subject' => '欢迎注册',
    ]);

    return ['message' => '注册成功,确认邮件稍后送达'];
  }
}

最佳实践

  1. 任务数据不宜过大:建议控制在 2MB 以内(与服务配置 OPTION_PACKAGE_MAX_LENGTH 一致),大内容只传标识符,由处理器自行从数据库、缓存或文件系统加载。
  2. 同步等待慎用emitWait() 会阻塞当前协程直到结果返回或超时,适合轻量查询型任务;耗时任务用 emit() 投递后通过缓存、数据库或其他渠道通知结果。
  3. 超时时间按任务实际耗时设置emitWait() 的默认 0.5 秒仅适合快速任务,慢任务请显式传入更大的超时值。
  4. 失败表达结构化:回传 ['ok' => false, 'error' => '...'],避免裸 false / null 与超时信号混淆。

下一步

  • 生命周期钩子:了解任务的 task / finish 事件在哪个进程触发,以及 Worker 启动钩子的用法
  • 服务提供者:将任务注册组织到模块的服务提供者中
  • 缓存:为任务队列配置高性能的 Redis 通道