异步任务管理
Viswoole 的异步任务系统基于 Swoole Task 机制封装,提供轻量级但功能完整的任务管理能力,支持任务注册、投递、队列持久化和结果等待。
核心组件
TaskManager - 任务管理器
负责任务的注册、投递和分发处理。
TaskProxy - 任务代理
封装 Swoole 原生 Task 对象,提供类型安全的任务数据访问,并通过 __get/__call 代理访问 SwooleTask 的属性与方法(如 finish())。
Task 门面
提供静态调用入口:
use Viswoole\Core\Facade\Task;前置配置
使用任务系统前,必须在 config/server.php 中启用 Task 相关选项:
// config/server.php
use Swoole\Constant;
return [
'servers' => [
'http' => [
'options' => [
// 任务进程数量,0 代表不开启任务能力
Constant::OPTION_TASK_WORKER_NUM => swoole_cpu_num(),
// 任务协程(框架在合并服务配置时会强制开启,此处可不配置)
Constant::OPTION_TASK_ENABLE_COROUTINE => true,
]
]
],
];注意: 使用任务系统必须满足
OPTION_TASK_USE_OBJECT => true或OPTION_TASK_ENABLE_COROUTINE => true至少开启一个。框架在Server合并服务配置时会强制开启OPTION_TASK_ENABLE_COROUTINE,因此该条件默认已满足,只需确保OPTION_TASK_WORKER_NUM > 0即可。
任务注册
警告: 请勿自行注册 Swoole 的
onTask事件(无论是通过ServerEventHook::addEvent('task', ...)还是$server->on('task', ...))。
- Swoole 层面:
Server->on()对同一事件重复注册时会覆盖前一次设置(官方文档明确说明),直接调用会完全替换掉框架的任务分发器,导致任务系统整体失效- 框架层面:
ServerEventHook虽支持同一事件挂载多个处理器,但 Swoole 所有服务端事件中仅onTask的返回值有语义(将任务结果回传给 Worker 进程并触发onFinish),多处理器时仅最后一个处理器的返回值生效——额外注册的处理器会导致任务被重复执行、emitWait()等同步等待拿到的结果被破坏正确做法:任务业务逻辑一律通过
Task::register()注册任务主题实现。
注册单个任务(回调方式)
use Viswoole\Core\Facade\Task;
use Viswoole\Core\Server\TaskProxy;
// 在服务启动前注册(如 bootstrap 或 Provider 中)
Task::register('sendEmail', function (TaskProxy $task, $server) {
$data = $task->data;
// 执行发送邮件逻辑...
$result = Mailer::send($data['to'], $data['subject'], $data['body']);
// 标记任务完成并返回结果
$task->finish([
'success' => true,
'message_id' => $result->getId(),
]);
});注册类批量任务
传入类名时,框架会通过反射自动扫描该类的全部方法(魔术方法除外),以 “主题.方法名” 格式批量注册:
class SmsService
{
/**
* 发送登录验证码
*/
public static function sendLoginCode(TaskProxy $task): void
{
$phone = $task->data['phone'];
$code = $task->data['code'];
SmsSender::send($phone, "您的验证码是:{$code},5分钟内有效。");
$task->finish(['status' => 'ok']);
}
/**
* 发送注册验证码
*/
public static function sendRegisterCode(TaskProxy $task): void
{
// ...
$task->finish(['status' => 'ok']);
}
/**
* 发送营销短信
*/
public function sendMarketing(TaskProxy $task): void
{
// 实例方法同样支持
$task->finish(['status' => 'ok']);
}
}
// 注册后自动生成以下主题:
// sms.sendLoginCode
// sms.sendRegisterCode
// sms.sendMarketing
Task::register('sms', SmsService::class);注册规则:
- 静态方法:注册为
类名前缀.方法名,回调为ClassName::methodName - 实例方法:注册为
类名前缀.方法名,回调为[ClassName, 'methodName'] - 以
__开头的魔术方法会被过滤
任务投递
emit() - 异步投递(非阻塞)
最常用的任务投递方式,立即返回不等待结果:
use Viswoole\Core\Facade\Task;
// 基本投递
$taskId = Task::emit('sendEmail', [
'to' => 'user@example.com',
'subject' => '欢迎注册',
'body' => '感谢您注册我们的服务...',
]);
if ($taskId !== false) {
echo "任务已投递,ID: {$taskId}";
}
// 不持久化到队列(服务重启后丢失)
Task::emit('cleanup', ['type' => 'temp'], queue: false);参数说明:
| 参数 | 类型 | 说明 |
|---|---|---|
$topic | string | 已注册的任务主题名称 |
$data | mixed | 传递给处理器的业务数据 |
$queue | bool | 是否持久化到队列(默认 true) |
返回值:
- 成功: 返回任务 ID (
int) - 失败: 返回
false
emitWait() - 同步等待(阻塞)
投递任务并阻塞等待执行结果。处理器通过 finish() 或 return 回传的数据类型原样透传(数组、整型等任意值):
try {
$result = Task::emitWait('generateReport', [
'user_id' => 100,
'type' => 'monthly',
], timeout: 5.0);
if ($result !== false) {
// $result 是处理器回传的原始数据,可能是数组、字符串等任意类型
print_r($result);
} else {
echo "任务超时、投递失败或处理器回传了 false";
}
} catch (\InvalidArgumentException $e) {
echo "任务主题不存在: " . $e->getMessage();
}重要:以下为 Swoole
taskwait的实测语义(Swoole 6.2 验证),编写任务处理器时务必注意:
处理器行为 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() - 并发多任务等待
同时投递多个任务并等待全部完成:
// 普通并发模式(taskWaitMulti)
$results = Task::emitsWait([
'getUserInfo' => ['user_id' => 100],
'getOrderList' => ['user_id' => 100],
], timeout: 3.0);
print_r($results);
// ['getUserInfo' => [...], 'getOrderList' => [...]]
// 协程并发模式(taskCo),更高效
$results = Task::emitsWait([
'task1' => [...],
'task2' => [...],
], timeout: 3.0, isCo: true);TaskProxy 使用详解
访问任务数据
Task::register('processData', function (TaskProxy $task) {
// 获取业务数据
$data = $task->data; // mixed - 投递时传入的数据
// 获取元信息
$topic = $task->topic; // string - 任务主题
$queueId = $task->queue_id; // string|null - 队列ID(非队列任务为null)
// 通过 __get 代理访问 SwooleTask 属性
$taskId = $task->id; // int - 任务ID
$workerId = $task->worker_id; // int - 所在Worker进程ID
$dispatchTime = $task->dispatch_time; // float - 投递时间戳
$flags = $task->flags; // int - 标志位
});回传任务结果
任务结果有两种回传方式,均遵循 Swoole 官方语义:
方式一:处理器 return——返回值经框架原样透传给 Worker 进程(触发 onFinish 事件 / taskwait 返回)。返回 null(含 void 处理器)不会触发结果投递:
Task::register('queryUser', function (TaskProxy $task) {
return ['id' => $task->data['id'], 'name' => '...']; // 直接 return 结果
});方式二:调用 finish()——通过 __call 代理直接透传给底层 SwooleTask,可以在任务处理器中多次调用,向 Worker 进程发送多个处理结果(每次都会触发 Worker 进程的 onFinish 事件):
Task::register('heavyComputation', function (TaskProxy $task) {
$input = $task->data['number'];
// 执行耗时计算...
$result = $this->compute($input);
// 调用 finish 向 Worker 进程返回结果
$task->finish([
'input' => $input,
'output' => $result,
'duration' => microtime(true) - $task->dispatch_time,
]);
// 流式任务可以多次调用 finish,逐步回传中间结果
// $task->finish(['progress' => 50]);
// $task->finish(['progress' => 100]);
});注意:
- 队列缓存清理与
finish()无关——队列任务在被 Task Worker 消费时即从缓存队列移除(见下文"队列持久化")。finish()只负责结果回传- 任务已过期或处理器抛出异常时,框架会向 Worker 进程返回
false(不会执行处理器 / 异常被记录到任务日志)
数据序列化
// 序列化(用于跨进程传输大数据)
$packed = TaskProxy::pack($largeData);
// 反序列化
$data = TaskProxy::unpack($packed);队列缓存通道
任务队列的持久化默认跟随 cache.default 指定的默认缓存通道。如果默认通道是文件缓存,而任务吞吐量较大,可以为任务队列单独指定高性能缓存通道(如 Redis),与业务缓存互不影响。
在 config/cache.php 中注册所需的缓存商店,然后在 config/task.php 中指定:
// config/cache.php
return [
'default' => env('cache.store', 'file'),
'stores' => [
'file' => Cache::FILE_DRIVER,
'redis' => ['driver' => Cache::REDIS_DRIVER, 'options' => [env('REDIS_HOST', '127.0.0.1')]],
],
];
// config/task.php
return [
// 任务队列持久化使用的缓存商店名称,需在 cache.stores 中已定义
// null 表示跟随 cache.default 默认通道
'store' => env('task.store', 'redis'),
];指定后,任务队列的标签记账、任务数据读写全部走该通道;业务缓存的默认通道不受影响。
注意:
task.store指向的商店必须已在cache.stores中定义,否则首次使用任务队列时会抛出CacheErrorException(快速失败)- 切换队列缓存通道前,残留在旧通道中的队列条目不会被新通道恢复。切换时请确保队列已清空(所有任务已消费)
任务有效期
部分任务具有时效性——比如验证码短信只在几分钟内有效,服务重启后恢复的队列任务如果已过时效,再执行只会浪费资源甚至产生副作用(给用户发送早已过期的验证码)。任务有效期用于在消费时淘汰这类过期任务。
在 config/task.php 中配置:
// config/task.php
return [
// 全局默认有效期(秒),从投递时刻起算;0 或 null 表示长期有效
'expire' => env('task.expire'),
// 按主题覆盖有效期(秒),优先级高于全局默认,主题名不区分大小写
// 值为 null 表示该主题长期有效(可用于全局设置了过期时间时为特定主题豁免)
'topics' => [
'sms.sendLoginCode' => 300, // 验证码任务 5 分钟有效
'report.generate' => 3600, // 报表任务 1 小时有效
'log.write' => null, // 日志任务永不过期
],
];生效机制:
- 有效期在
emit()投递时按主题解析(task.topics>task.expire),并随任务数据一起传递(含持久化到队列的数据) - 任务被 Task Worker 消费时检测:已过期的任务不会执行处理器,直接向 Worker 进程返回
false(onFinish回调收到false,taskwait立即返回false),并记录一条任务日志 - 仅对
emit()投递的任务生效;emitWait()/emitsWait()为同步执行,不受此配置影响
提示:过期跳过的任务同样遵循队列的 at-most-once 语义——消费时即从队列删除,不会重投。
队列持久化
当 emit() 的 $queue 参数为 true(默认值)时,任务会被持久化到缓存队列中:
// 投递任务并持久化
$taskId = Task::emit('importantJob', $data);
// 此时任务信息已写入缓存,包含:
// - taskData(原始数据)
// - topic(主题)
// - queueId(唯一标识)
// - dispatchTime(投递时间)队列语义(at-most-once):
任务队列是待消费队列——任务被 Task Worker 消费时立即从缓存队列中删除,与执行结果无关(finish() 仅向 Worker 进程回传结果,框架无法据此判定任务成败)。
因此崩溃恢复只覆盖一个窗口:任务已投递(写入缓存)但尚未被 Task Worker 消费。服务重启时框架会自动重新投递这一部分任务:
- Worker 进程在任务消费前崩溃(如
kill -9)→ 队列条目仍在缓存中 → 重启后自动恢复重跑 - 任务已被消费(无论执行成功、失败还是抛异常)→ 条目已删除 → 不会重投
注意:这意味着任务执行中途进程崩溃不会重试(at-most-once 语义)。如果业务需要"执行失败重试",应在任务处理器内部自行实现(如捕获异常后重新
emit,或回传失败结果由调用方决定)。
服务重启恢复:
框架会在 WorkerStart 事件中自动检测未消费的队列任务并重新投递:
// 内部实现(无需手动编写)
ServerEventHook::addEvent('workerStart', function ($server, $workerId) {
if (!$server->taskworker) {
$store = $this->getQueueCacheStore((string)$workerId);
$taskQueue = $store->get();
foreach ($taskQueue as $queueId) {
go(function () use ($queueId, $store) {
$taskData = $this->cache->get($queueId);
if ($taskData) {
Server::getServer()->task($taskData);
}
});
}
}
});完整示例
异步邮件发送服务
namespace App\Service;
use Viswoole\Core\Server\TaskProxy;
use Viswoole\Core\Facade\Task;
class EmailTaskService
{
/**
* 注册邮件相关任务
*/
public static function register(): void
{
Task::register('email.send', [self::class, 'handleSend']);
Task::register('email.batch', [self::class, 'handleBatch']);
}
/**
* 发送单封邮件
*/
public static function handleSend(TaskProxy $task): void
{
$to = $task->data['to'];
$subject = $task->data['subject'];
$body = $task->data['body'];
try {
$messageId = app(Mailer::class)->send($to, $subject, $body);
$task->finish([
'success' => true,
'message_id' => $messageId,
]);
} catch (\Throwable $e) {
$task->finish([
'success' => false,
'error' => $e->getMessage(),
]);
}
}
/**
* 批量发送邮件
*/
public static function handleBatch(TaskProxy $task): void
{
$recipients = $task->data['recipients']; // [['to'=>...,'subject'=>...,'body'=>...]]
$template = $task->data['template'];
$results = [];
foreach ($recipients as $recipient) {
try {
$id = app(Mailer::class)->send(
$recipient['to'],
$recipient['subject'] ?? $template['subject'],
$recipient['body'] ?? $template['body']
);
$results[] = ['to' => $recipient['to'], 'success' => true, 'id' => $id];
} catch (\Throwable $e) {
$results[] = ['to' => $recipient['to'], 'success' => false, 'error' => $e->getMessage()];
}
}
$task->finish(['results' => $results, 'total' => count($results)]);
}
}
// 在控制器中使用
class UserController extends Controller
{
public function register(Request $request)
{
// 验证数据...
// 创建用户
$user = User::create($request->validated());
// 异步发送欢迎邮件(不阻塞响应)
Task::emit('email.send', [
'to' => $user->email,
'subject' => '欢迎加入',
'body' => view('emails.welcome', ['user' => $user]),
]);
return json(['code' => 0, 'msg' => '注册成功']);
}
}最佳实践
1. 错误处理
框架在 onTask 分发时已对 Throwable 做兜底:处理器抛出的异常会被捕获并记录到任务日志通道(含主题、队列ID、异常原文和任务数据),同时向 Worker 进程返回 false(taskwait 调用方立即收到 false 而非干等超时),不会导致 Task Worker 进程崩溃。
业务层面仍建议自行捕获异常并回传结构化结果(false 无法与超时/过期区分):
Task::register('riskyOperation', function (TaskProxy $task) {
try {
$result = doSomethingRisky($task->data);
$task->finish(['ok' => true, 'data' => $result]);
} catch (\Throwable $e) {
// 配合 emitWait 使用时,务必回传结构化结果而非裸 false/null:
// - 裸 false 与超时返回值无法区分
// - 裸 null / 无返回值会导致 emitWait 调用方阻塞满整个超时时间
$task->finish([
'ok' => false,
'error' => $e->getMessage(),
'trace' => $e->getTraceAsString(),
]);
}
});2. 超时控制
// 设置合理的超时时间
$result = Task::emitWait('slowTask', $data, timeout: 30.0);
// 对于可能长时间运行的任务,建议使用异步模式 + 回调通知
Task::emit('longRunningTask', $data);
// 任务完成后通过其他渠道(如 WebSocket、轮询 API)通知前端3. 数据大小限制
// Task 数据不宜过大(建议 < 2MB),大数据请传递标识符
Task::emit('processFile', [
'file_id' => $fileId, // 传ID而非文件内容
'user_id' => $userId,
]);
// 在任务处理器中根据 ID 从数据库/Redis/文件系统获取实际数据