PHP 怎么TaskWorker

wen PHP项目 2

本文目录导读:

PHP 怎么TaskWorker

  1. 使用 Swoole 的 TaskWorker(推荐)
  2. 使用 Laravel/ThinkPHP 的队列系统
  3. 使用进程管理扩展(pcntl + posix)
  4. 使用消息队列(Redis + PHP)
  5. 使用常驻内存框架
  6. 最佳实践建议

在 PHP 中实现 TaskWorker 模式(任务工作者模式),主要有以下几种方案:

使用 Swoole 的 TaskWorker(推荐)

Swoole 内置了 TaskWorker 机制,是最成熟的解决方案:

<?php
$server = new Swoole\Server('0.0.0.0', 9501);
// 设置任务进程数量
$server->set([
    'worker_num' => 4,        // Worker 进程数
    'task_worker_num' => 8,   // TaskWorker 进程数
    'task_enable_coroutine' => true,  // 启用协程支持
]);
// 处理异步任务
$server->on('Task', function ($server, $task_id, $from_id, $data) {
    echo "接收任务: $task_id\n";
    // 执行耗时操作
    $result = processTask($data);
    // 返回结果给 Worker
    return $result;
});
// 处理任务完成
$server->on('Finish', function ($server, $task_id, $data) {
    echo "任务完成: $task_id, 结果: $data\n";
});
// 投递任务
$server->on('Receive', function ($server, $fd, $reactor_id, $data) {
    // 投递任务到 TaskWorker
    $task_id = $server->task($data);
    echo "任务已投递: $task_id\n";
});
$server->start();
function processTask($data) {
    // 模拟耗时操作
    sleep(2);
    return "处理结果: " . json_encode($data);
}

使用协程版 TaskWorker

<?php
use Swoole\Coroutine;
$server = new Swoole\Http\Server('0.0.0.0', 9502);
$server->set([
    'worker_num' => 2,
    'task_worker_num' => 4,
    'task_enable_coroutine' => true,
]);
$server->on('Request', function ($request, $response) use ($server) {
    // 投递异步任务
    $result = $server->taskCo([
        ['type' => 'email', 'data' => '发送邮件'],
        ['type' => 'log', 'data' => '写入日志'],
        ['type' => 'report', 'data' => '生成报表'],
    ], 10);  // 10秒超时
    $response->end("任务完成: " . json_encode($result));
});
// 使用协程处理任务
$server->on('Task', function ($server, $task) {
    if (isset($task->data['type'])) {
        switch ($task->data['type']) {
            case 'email':
                co::sleep(2);  // 协程等待
                $task->finish(['status' => 'email_sent']);
                break;
            case 'log':
                co::sleep(1);
                $task->finish(['status' => 'log_written']);
                break;
            default:
                $task->finish(['status' => 'unknown_type']);
        }
    }
});
$server->start();

使用 Laravel/ThinkPHP 的队列系统

Laravel 队列 Worker

// Laravel 配置文件 config/queue.php
'connections' => [
    'redis' => [
        'driver' => 'redis',
        'connection' => 'default',
        'queue' => env('REDIS_QUEUE', 'default'),
        'retry_after' => 90,
        'block_for' => 0,
    ],
],
// 创建任务类
namespace App\Jobs;
class ProcessPodcast implements ShouldQueue
{
    use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
    protected $podcast;
    public function __construct($podcast)
    {
        $this->podcast = $podcast;
    }
    public function handle()
    {
        // 处理任务
        echo "处理播客: {$this->podcast['name']}\n";
        sleep(2);  // 模拟耗时
        return true;
    }
}
// 分发任务
ProcessPodcast::dispatch(['name' => 'Test', 'duration' => 120]);
// 启动队列 Worker
// php artisan queue:work --daemon --tries=3
// 或
// php artisan queue:listener --timeout=60

多进程队列处理

# 使用 Laravel 的 queue:work 多进程
php artisan queue:work --queue=high,default --timeout=60 --sleep=15 --tries=3
# 使用 Supervisor 管理多个 worker

使用进程管理扩展(pcntl + posix)

<?php
class TaskWorker {
    private $workers = [];
    private $taskQueue = [];
    private $maxWorkers = 4;
    public function __construct($maxWorkers = 4) {
        $this->maxWorkers = $maxWorkers;
        $this->initProcessPool();
    }
    // 初始化进程池
    private function initProcessPool() {
        for ($i = 0; $i < $this->maxWorkers; $i++) {
            $pid = pcntl_fork();
            if ($pid == -1) {
                die("无法创建子进程\n");
            } elseif ($pid) {
                // 父进程
                $this->workers[$pid] = $i;
            } else {
                // 子进程
                $this->workerProcess($i);
                exit(0);
            }
        }
    }
    // 子进程处理函数
    private function workerProcess($workerId) {
        // 创建消息队列
        $msgKey = ftok(__FILE__, 's');
        $msgQueue = msg_get_queue($msgKey, 0666);
        echo "Worker #{$workerId} 已启动, PID: " . getmypid() . "\n";
        while (true) {
            // 接收任务
            if (msg_receive($msgQueue, 1, $msgType, 1024, $message, true)) {
                echo "Worker #{$workerId} 处理任务: {$message}\n";
                // 模拟处理耗时任务
                sleep(2);
                // 记录处理结果
                $result = "任务 [{$message}] 由 Worker #{$workerId} 处理完成\n";
                file_put_contents('/tmp/task_result.log', $result, FILE_APPEND);
            }
            // 检查是否有停止信号
            pcntl_signal_dispatch();
        }
    }
    // 投递任务
    public function dispatch($task) {
        $msgKey = ftok(__FILE__, 's');
        $msgQueue = msg_get_queue($msgKey, 0666);
        // 发送消息到队列
        msg_send($msgQueue, 1, $task, true);
        echo "任务已投递: {$task}\n";
    }
    // 监控进程
    public function monitor() {
        while (true) {
            $status = 0;
            $pid = pcntl_wait($status, WNOHANG);
            if ($pid > 0) {
                echo "Worker {$pid} 已退出,状态: {$status}\n";
                // 重新创建工人进程
                $newPid = pcntl_fork();
                if ($newPid == 0) {
                    $this->workerProcess($this->workers[$pid]);
                }
            }
            sleep(1);
        }
    }
}
// 使用示例
$taskWorker = new TaskWorker(4);
$taskWorker->dispatch('发送邮件');
$taskWorker->dispatch('生成报表');
$taskWorker->dispatch('处理图片');
$taskWorker->monitor();

使用消息队列(Redis + PHP)

<?php
// Redis 任务队列 Worker
class RedisQueueWorker {
    private $redis = null;
    private $queueName = 'task_queue';
    public function __construct() {
        $this->redis = new Redis();
        $this->redis->connect('127.0.0.1', 6379);
    }
    // 投递任务
    public function dispatch($task) {
        $taskData = json_encode([
            'id' => uniqid('task_'),
            'data' => $task,
            'created_at' => time(),
        ]);
        return $this->redis->lPush($this->queueName, $taskData);
    }
    // 处理任务
    public function run() {
        echo "Worker 已启动,等待任务... \n";
        while (true) {
            // 阻塞获取任务
            $task = $this->redis->brPop([$this->queueName], 30);
            if ($task) {
                $taskData = json_decode($task[1], true);
                echo "处理任务: {$taskData['id']}\n";
                echo "任务数据: " . json_encode($taskData['data']) . "\n";
                // 模拟耗时任务
                sleep(2);
                echo "任务处理完成\n\n";
            }
        }
    }
}
// 使用示例
$worker = new RedisQueueWorker();
// 生产:投递任务
$worker->dispatch(['type' => 'email', 'to' => 'user@example.com']);
$worker->dispatch(['type' => 'report', 'date' => '2024-01-01']);
// 消费:运行多个 worker
// 可以开启多个终端运行脚本文件 worker.php
//$worker->run();

使用常驻内存框架

Workerman

<?php
require_once __DIR__ . '/vendor/autoload.php';
use Workerman\Worker;
use Workerman\Timer;
// 创建消息队列
$taskWorker = new Worker('tcp://0.0.0.0:1234');
$taskWorker->count = 8;  // 设置8个进程
// 收到任务
$taskWorker->onMessage = function ($connection, $task) {
    echo "收到任务: {$task}\n";
    // 模拟耗时任务
    sleep(2);
    // 返回结果
    $connection->send("任务处理完成: {$task}");
};
// 创建一个任务投递 Worker
$taskProducer = new Worker('text://0.0.0.0:1235');
$taskProducer->count = 1;
// 投递任务到 TaskWorker
$taskProducer->onMessage = function ($connection, $data) {
    // 连接 TaskWorker
    $client = stream_socket_client('tcp://127.0.0.1:1234');
    fwrite($client, $data);
    // 获取结果
    $result = fread($client, 1024);
    fclose($client);
    $connection->send("结果: {$result}");
};
Worker::runAll();

最佳实践建议

  1. 选择依据

    • 已有 Swoole 项目 → 使用 Swoole TaskWorker
    • Laravel 项目 → 使用队列系统
    • 轻量级项目 → 使用 Redis 队列 + 简单 Worker
    • 需要复杂进程管理 → Workerman/Swoole
  2. 性能优化

    • 设置合适的 Worker 数量(通常为 CPU 核心数)
    • 使用协程处理 I/O 密集任务
    • 合理设置超时和重试机制
  3. 监控与维护

    • 使用 Supervisor 管理进程
    • 记录任务日志
    • 监控队列长度和 Worker 状态
  4. 注意事项

    • 避免在 TaskWorker 中执行阻塞操作
    • 任务尽量幂等(可重复执行)
    • 确保任务失败有重试机制
    • 合理设置内存限制防止内存泄漏

选择哪种方案取决于你的项目架构、性能需求和团队熟悉程度。

抱歉,评论功能暂时关闭!