PHP队列系统与异步任务处理

队列是处理异步任务的标准方式。邮件发送、图片处理、报表生成这些耗时操作放到队列中异步处理,可以提升用户体验。今天说说PHP队列的实现。

用Redis的List实现简单队列。

```php
class RedisQueue
{
private Redis $redis;
private string $prefix = 'queue:';

public function __construct(Redis $redis)
{
$this->redis = $redis;
}

public function push(string $queue, array $job): string
{
$jobId = uniqid('job_', true);
$payload = json_encode(['id' => $jobId, 'data' => $job, 'created_at' => time(), 'attempts' => 0]);
$this->redis->lPush($this->prefix . $queue, $payload);
return $jobId;
}

public function pop(string $queue, int $timeout = 0): ?array
{
$result = $this->redis->brPop([$this->prefix . $queue], $timeout);
if ($result === null) return null;
$payload = json_decode($result[1], true);
if ($payload === null) return null;
$payload['attempts']++;
return $payload;
}

public function later(string $queue, array $job, int $delay): void
{
$payload = json_encode(['id' => uniqid('job_', true), 'data' => $job, 'created_at' => time(), 'available_at' => time() + $delay]);
$this->redis->zAdd($this->prefix . 'delayed:' . $queue, time() + $delay, $payload);
}

public function size(string $queue): int
{
return $this->redis->lLen($this->prefix . $queue);
}

public function clear(string $queue): void
{
$this->redis->del($this->prefix . $queue);
}
}
?>

Worker进程从队列中拉取并执行任务。

```php
class QueueWorker
{
private RedisQueue $queue;
private array $handlers = [];
private bool $running = false;

public function __construct(RedisQueue $queue)
{
$this->queue = $queue;
}

public function addHandler(string $type, callable $handler): void
{
$this->handlers[$type] = $handler;
}

public function work(string $queueName, int $memoryLimit = 128): void
{
$this->running = true;
echo "Worker启动,监听: $queueName\n";

while ($this->running) {
if (memory_get_usage(true) > $memoryLimit * 1024 * 1024) {
echo "内存超限,退出\n";
break;
}

try {
$job = $this->queue->pop($queueName, 5);
if ($job === null) continue;
$this->processJob($job);
} catch (Exception $e) {
echo "错误: {$e->getMessage()}\n";
sleep(1);
}
}
}

private function processJob(array $job): void
{
$type = $job['data']['type'] ?? 'default';
$handler = $this->handlers[$type] ?? null;

if ($handler === null) {
echo "无处理器: $type\n";
return;
}

try {
$handler($job['data']);
echo "任务完成: {$job['id']}\n";
} catch (Exception $e) {
echo "任务失败: {$job['id']}: {$e->getMessage()}\n";
}
}

public function stop(): void
{
$this->running = false;
}
}

$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$queue = new RedisQueue($redis);
$worker = new QueueWorker($queue);

$worker->addHandler('send_email', function ($data) {
echo "发送邮件至 {$data['to']}\n";
sleep(1);
});

$worker->addHandler('resize_image', function ($data) {
echo "调整图片: {$data['path']}\n";
sleep(1);
});

// 推送任务
$queue->push('default', ['type' => 'send_email', 'to' => 'user@example.com', 'subject' => '欢迎']);
$queue->push('default', ['type' => 'resize_image', 'path' => '/uploads/photo.jpg', 'sizes' => [300, 600]]);

echo "队列大小: " . $queue->size('default') . "\n";
// 生产环境在CLI中启动Worker
// $worker->work('default');
?>

队列系统可以解决很多实际问题。异步任务处理、削峰填谷、任务重试、定时执行。用Redis做队列比用消息队列系统简单得多,功能也够用。生产环境建议用supervisor管理Worker进程,确保进程退出后自动重启。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐