异步任务
异步任务(Async Task)用于把耗时代码——发送邮件、生成报表、推送消息等——从请求处理流程中剥离,交由独立的 Task Worker 进程执行,避免阻塞 Worker 响应请求。Viswoole 基于 Swoole Task 机制封装了任务管理器(TaskManager),提供主题注册、异步投递、队列持久化、有效期淘汰与同步等待结果等能力。
核心组件
| 组件 | 全限定名 | 职责 |
|---|---|---|
| TaskManager | Viswoole\Core\Server\TaskManager | 任务主题的注册、投递与分发 |
| TaskProxy | Viswoole\Core\Server\TaskProxy | 任务代理对象,作为处理器入参 |
| Task 门面 | Viswoole\Core\Facade\Task | 静态代理 TaskManager 的核心方法 |
| TaskService | Viswoole\Core\Service\TaskService | 任务服务提供者,负责装载任务管理器 |
TaskService 已在框架默认 config/app.php 的 services[] 中注册,任务能力开箱即用。业务代码统一通过 Task 门面调用:
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 核数开启任务进程,一般无需修改:
// 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)」组织:投递时指定主题名,框架将任务分发给该主题注册的处理器。主题名不区分大小写。
注册单个主题(回调方式)
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() 的第二个参数传入类名时,框架通过反射扫描该类的全部公开方法,以 「主题前缀.方法名」 的格式批量注册:
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'];
}
}// 注册后生成 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() 投递任务后立即返回,不等待执行结果,是最常用的投递方式:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
$topic | string | 无 | 已注册的任务主题名称,未注册时抛出 InvalidArgumentException |
$data | mixed | 无 | 传递给处理器的业务数据 |
$queue | bool | true | 是否持久化到队列,启用后服务重启可自动恢复未消费的任务 |
返回值:投递成功返回任务 ID(int),投递失败返回 false。$queue = true 且队列缓存写入失败时抛出 RuntimeException。
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 中配置:
| 配置键 | 类型 | 默认值 | 说明 |
|---|---|---|---|
store | string|null | null | 任务队列持久化使用的缓存商店名称,需在 config/cache.php 的 stores 中定义;null 表示跟随 cache.default 默认通道 |
expire | int|null | null | 全局默认有效期(秒),从投递时刻起算;0 或 null 表示长期有效 |
topics | array<string, int|null> | [] | 按主题覆盖有效期(秒),优先级高于 expire,主题名不区分大小写;值为 null 表示该主题豁免(长期有效) |
// 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 回传的数据类型原样透传:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
$topic | string | 无 | 已注册的任务主题名称 |
$data | mixed | 无 | 传递给处理器的业务数据 |
$timeout | float | 0.5 | 等待超时时间(秒) |
返回值:处理器回传的原始结果;主题未注册时抛出 InvalidArgumentException;超时或投递失败返回 false。
$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()
同时投递多个任务并等待全部完成:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
$tasks | array<string, mixed> | 无 | 任务列表,键为主题名称,值为业务数据 |
$timeout | float | 0.5 | 等待超时时间(秒) |
$isCo | bool | false | true 使用协程并发调度(taskCo),false 使用 taskWaitMulti |
返回值:各任务的执行结果列表,超时未完成的任务对应位置为 false。
$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 服务实例,一般无需使用:
function (TaskProxy $task, \Swoole\Server $server): mixed {}TaskProxy 属性
| 属性 | 类型 | 默认值 | 说明 |
|---|---|---|---|
data | mixed | null | 投递时传入的业务数据 |
topic | string | 无 | 任务主题名称 |
queue_id | string|null | null | 队列唯一标识,非队列任务为 null |
除上述属性外,TaskProxy 还通过 __get 代理访问 Swoole 原生 Task 对象的属性,如 id(任务 ID)、worker_id(所在 Worker 进程 ID)、dispatch_time(投递时间戳)、flags(标志位);访问不存在的属性会抛出异常。
回传结果:return 与 finish()
方式一:处理器 return——返回值经框架原样透传给 Worker 进程(触发 onFinish 事件或 taskwait 返回)。返回 null(含 void 处理器)不会触发结果投递,异步任务(配合 emit())通常直接用这种方式:
Task::register('email.query', function (TaskProxy $task): array {
return ['id' => $task->data['id'], 'status' => 'sent'];
});方式二:调用 finish(mixed $data): bool——通过 __call 直接透传给底层 Swoole Task 对象,遵循 Swoole 语义可以在任务处理器中多次调用,向 Worker 进程发送多个处理结果(流式任务可逐步回传中间结果):
Task::register('report.build', function (TaskProxy $task): void {
$task->finish(['progress' => 50]); // 中间进度
// ...继续计算
$task->finish(['progress' => 100]); // 最终结果
});注意:队列条目的清理发生在任务被消费时,与 finish() 无关——finish() 只负责结果回传。
异常处理
框架在任务分发时已对 Throwable 做兜底:处理器抛出的异常会被捕获并记录到任务日志通道(含主题、队列 ID、异常原文和任务数据),同时向 Worker 进程返回 false(emitWait 调用方立即收到失败信号而不是干等超时),不会导致 Task Worker 进程崩溃。
业务层面仍建议自行捕获异常并回传结构化结果——裸 false 无法与超时/过期区分:
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 的静态方法手动序列化/反序列化:
$packed = TaskProxy::pack($largeData); // 序列化为二进制字符串,失败返回 false
$data = TaskProxy::unpack($packed); // 反序列化,失败返回 false完整示例
控制器中注册用户后异步发送欢迎邮件,不阻塞响应:
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' => '注册成功,确认邮件稍后送达'];
}
}最佳实践
- 任务数据不宜过大:建议控制在 2MB 以内(与服务配置
OPTION_PACKAGE_MAX_LENGTH一致),大内容只传标识符,由处理器自行从数据库、缓存或文件系统加载。 - 同步等待慎用:
emitWait()会阻塞当前协程直到结果返回或超时,适合轻量查询型任务;耗时任务用emit()投递后通过缓存、数据库或其他渠道通知结果。 - 超时时间按任务实际耗时设置:
emitWait()的默认 0.5 秒仅适合快速任务,慢任务请显式传入更大的超时值。 - 失败表达结构化:回传
['ok' => false, 'error' => '...'],避免裸false/null与超时信号混淆。
