本文目录导读:

- 使用 Swoole 的 TaskWorker(推荐)
- 使用 Laravel/ThinkPHP 的队列系统
- 使用进程管理扩展(pcntl + posix)
- 使用消息队列(Redis + PHP)
- 使用常驻内存框架
- 最佳实践建议
在 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();
最佳实践建议
-
选择依据:
- 已有 Swoole 项目 → 使用 Swoole TaskWorker
- Laravel 项目 → 使用队列系统
- 轻量级项目 → 使用 Redis 队列 + 简单 Worker
- 需要复杂进程管理 → Workerman/Swoole
-
性能优化:
- 设置合适的 Worker 数量(通常为 CPU 核心数)
- 使用协程处理 I/O 密集任务
- 合理设置超时和重试机制
-
监控与维护:
- 使用 Supervisor 管理进程
- 记录任务日志
- 监控队列长度和 Worker 状态
-
注意事项:
- 避免在 TaskWorker 中执行阻塞操作
- 任务尽量幂等(可重复执行)
- 确保任务失败有重试机制
- 合理设置内存限制防止内存泄漏
选择哪种方案取决于你的项目架构、性能需求和团队熟悉程度。