异步任务管理

Viswoole 的异步任务系统基于 Swoole Task 机制封装,提供轻量级但功能完整的任务管理能力,支持任务注册、投递、队列持久化和结果等待。

核心组件

TaskManager - 任务管理器

负责任务的注册、投递和分发处理。

TaskProxy - 任务代理

封装 Swoole 原生 Task 对象,提供类型安全的任务数据访问,并通过 __get/__call 代理访问 SwooleTask 的属性与方法(如 finish())。

Task 门面

提供静态调用入口:

php
use Viswoole\Core\Facade\Task;

前置配置

使用任务系统前,必须在 config/server.php 中启用 Task 相关选项:

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 => true,
            ]
        ]
    ],
];

注意: 使用任务系统必须满足 OPTION_TASK_USE_OBJECT => trueOPTION_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() 注册任务主题实现。

注册单个任务(回调方式)

php
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(),
    ]);
});

注册类批量任务

传入类名时,框架会通过反射自动扫描该类的全部方法(魔术方法除外),以 “主题.方法名” 格式批量注册:

php
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() - 异步投递(非阻塞)

最常用的任务投递方式,立即返回不等待结果:

php
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);

参数说明:

参数类型说明
$topicstring已注册的任务主题名称
$datamixed传递给处理器的业务数据
$queuebool是否持久化到队列(默认 true

返回值:

  • 成功: 返回任务 ID (int)
  • 失败: 返回 false

emitWait() - 同步等待(阻塞)

投递任务并阻塞等待执行结果。处理器通过 finish()return 回传的数据类型原样透传(数组、整型等任意值):

php
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() - 并发多任务等待

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

php
// 普通并发模式(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 使用详解

访问任务数据

php
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 处理器)不会触发结果投递:

php
Task::register('queryUser', function (TaskProxy $task) {
    return ['id' => $task->data['id'], 'name' => '...']; // 直接 return 结果
});

方式二:调用 finish()——通过 __call 代理直接透传给底层 SwooleTask,可以在任务处理器中多次调用,向 Worker 进程发送多个处理结果(每次都会触发 Worker 进程的 onFinish 事件):

php
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(不会执行处理器 / 异常被记录到任务日志)

数据序列化

php
// 序列化(用于跨进程传输大数据)
$packed = TaskProxy::pack($largeData);

// 反序列化
$data = TaskProxy::unpack($packed);

队列缓存通道

任务队列的持久化默认跟随 cache.default 指定的默认缓存通道。如果默认通道是文件缓存,而任务吞吐量较大,可以为任务队列单独指定高性能缓存通道(如 Redis),与业务缓存互不影响。

config/cache.php 中注册所需的缓存商店,然后在 config/task.php 中指定:

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 中配置:

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 进程返回 falseonFinish 回调收到 falsetaskwait 立即返回 false),并记录一条任务日志
  • 仅对 emit() 投递的任务生效;emitWait() / emitsWait() 为同步执行,不受此配置影响

提示:过期跳过的任务同样遵循队列的 at-most-once 语义——消费时即从队列删除,不会重投。

队列持久化

emit()$queue 参数为 true(默认值)时,任务会被持久化到缓存队列中:

php
// 投递任务并持久化
$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 事件中自动检测未消费的队列任务并重新投递:

php
// 内部实现(无需手动编写)
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);
                }
            });
        }
    }
});

完整示例

异步邮件发送服务

php
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 进程返回 falsetaskwait 调用方立即收到 false 而非干等超时),不会导致 Task Worker 进程崩溃。

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

php
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. 超时控制

php
// 设置合理的超时时间
$result = Task::emitWait('slowTask', $data, timeout: 30.0);

// 对于可能长时间运行的任务,建议使用异步模式 + 回调通知
Task::emit('longRunningTask', $data);
// 任务完成后通过其他渠道(如 WebSocket、轮询 API)通知前端

3. 数据大小限制

php
// Task 数据不宜过大(建议 < 2MB),大数据请传递标识符
Task::emit('processFile', [
    'file_id' => $fileId,      // 传ID而非文件内容
    'user_id' => $userId,
]);

// 在任务处理器中根据 ID 从数据库/Redis/文件系统获取实际数据