本文目录导读:

在PHP项目中实现优先级队列以优先处理重要任务,主要有以下几种方案:
使用 PHP 内置的 SplPriorityQueue
<?php
class TaskPriorityQueue extends SplPriorityQueue {
public function compare($priority1, $priority2) {
// 数值越大优先级越高
if ($priority1 === $priority2) return 0;
return $priority1 > $priority2 ? 1 : -1;
}
}
// 使用示例
$queue = new TaskPriorityQueue();
// 添加任务(数据,优先级)
$queue->insert(['task' => '发送邮件', 'data' => [...]], 1); // 低优先级
$queue->insert(['task' => '支付回调处理', 'data' => [...]], 10); // 高优先级
$queue->insert(['task' => '用户注册通知', 'data' => [...]], 5); // 中优先级
// 取出任务(优先级高的先取出)
while ($queue->count() > 0) {
$task = $queue->extract();
processTask($task);
}
Redis 实现优先级队列
使用 Redis 的有序集合(ZSET):
<?php
class PriorityQueue {
private $redis;
private $queueKey;
public function __construct($queueKey = 'task_queue') {
$this->redis = new Redis();
$this->redis->connect('127.0.0.1', 6379);
$this->queueKey = $queueKey;
}
// 添加任务(优先级越高,score越小)
public function enqueue($taskData, $priority = 5) {
$taskId = uniqid('task_', true);
$task = json_encode([
'id' => $taskId,
'data' => $taskData,
'created_at' => time()
]);
// score = 优先级值(1-10,1为最高优先级)
$this->redis->zAdd($this->queueKey, $priority, $task);
}
// 获取最高优先级任务
public function dequeue() {
// 获取score最小的(优先级最高)
$tasks = $this->redis->zRange($this->queueKey, 0, 0);
if (!empty($tasks)) {
$task = $tasks[0];
// 从队列中移除
$this->redis->zRem($this->queueKey, $task);
return json_decode($task, true);
}
return null;
}
// 批量获取多个任务
public function dequeueBatch($limit = 10) {
$tasks = $this->redis->zRange($this->queueKey, 0, $limit - 1);
$result = [];
foreach ($tasks as $task) {
$this->redis->zRem($this->queueKey, $task);
$result[] = json_decode($task, true);
}
return $result;
}
}
// 使用示例
$queue = new PriorityQueue();
// 添加不同优先级的任务
$queue->enqueue(['type' => 'email', 'to' => 'user@example.com'], 5); // 普通
$queue->enqueue(['type' => 'payment', 'order_id' => 123], 1); // 高优先级
$queue->enqueue(['type' => 'log', 'message' => 'test'], 10); // 低优先级
// 处理任务(总是先处理高优先级)
while ($task = $queue->dequeue()) {
processTask($task);
}
RabbitMQ 实现优先级队列
<?php
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
// 生产者
class PriorityQueueProducer {
private $connection;
private $channel;
public function __construct() {
$this->connection = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$this->channel = $this->connection->channel();
// 声明优先级队列(需要x-max-priority参数)
$args = new \PhpAmqpLib\Wire\AMQPTable([
'x-max-priority' => 10 // 最大优先级10
]);
$this->channel->queue_declare('priority_tasks', false, true, false, false, false, $args);
}
public function publish($message, $priority = 0) {
$msg = new AMQPMessage(json_encode($message), [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
'priority' => $priority
]);
$this->channel->basic_publish($msg, '', 'priority_tasks');
}
public function __destruct() {
$this->channel->close();
$this->connection->close();
}
}
// 消费者
class PriorityQueueConsumer {
private $connection;
private $channel;
public function __construct() {
$this->connection = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$this->channel = $this->connection->channel();
// 每次只取1个任务(确保高优先级先处理)
$this->channel->basic_qos(null, 1, null);
}
public function consume() {
$callback = function ($msg) {
$task = json_decode($msg->body, true);
echo "处理任务: " . json_encode($task) . "\n";
echo "优先级: " . $msg->getPriority() . "\n";
// 处理任务
$this->processTask($task);
// 确认消息
$msg->ack();
};
$this->channel->basic_consume('priority_tasks', '', false, false, false, false, $callback);
while ($this->channel->is_consuming()) {
$this->channel->wait();
}
}
private function processTask($task) {
// 实际任务处理逻辑
sleep(1); // 模拟处理时间
}
}
多队列实现(简单但有效)
<?php
class MultiQueuePriority {
private $queues = [];
public function __construct() {
// 创建3个优先级队列
$this->queues = [
'high' => new SplQueue(),
'medium' => new SplQueue(),
'low' => new SplQueue()
];
}
public function addTask($task, $priority = 'medium') {
switch ($priority) {
case 'high':
$this->queues['high']->enqueue($task);
break;
case 'medium':
$this->queues['medium']->enqueue($task);
break;
case 'low':
$this->queues['low']->enqueue($task);
break;
default:
$this->queues['medium']->enqueue($task);
}
}
public function getNextTask() {
// 优先检查高优先级队列
if (!$this->queues['high']->isEmpty()) {
return ['priority' => 'high', 'task' => $this->queues['high']->dequeue()];
}
// 然后检查中优先级
if (!$this->queues['medium']->isEmpty()) {
return ['priority' => 'medium', 'task' => $this->queues['medium']->dequeue()];
}
// 最后检查低优先级
if (!$this->queues['low']->isEmpty()) {
return ['priority' => 'low', 'task' => $this->queues['low']->dequeue()];
}
return null;
}
// 带权重的获取方式(高优先级任务更频繁)
public function getNextTaskWeighted() {
$rand = mt_rand(1, 100);
// 70%概率处理高优先级
if ($rand <= 70 && !$this->queues['high']->isEmpty()) {
return $this->queues['high']->dequeue();
}
// 20%概率处理中优先级
elseif ($rand <= 90 && !$this->queues['medium']->isEmpty()) {
return $this->queues['medium']->dequeue();
}
// 10%概率处理低优先级
elseif (!$this->queues['low']->isEmpty()) {
return $this->queues['low']->dequeue();
}
// 回退到非空队列
elseif (!$this->queues['high']->isEmpty()) {
return $this->queues['high']->dequeue();
}
elseif (!$this->queues['medium']->isEmpty()) {
return $this->queues['medium']->dequeue();
}
elseif (!$this->queues['low']->isEmpty()) {
return $this->queues['low']->dequeue();
}
return null;
}
}
选择建议
| 方案 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
| SplPriorityQueue | 单进程内存队列 | 简单、内置 | 不支持持久化、跨进程 |
| Redis ZSET | 分布式系统 | 持久化、支持分布式 | 需要维护Redis |
| RabbitMQ | 企业级消息系统 | 功能完善、可靠 | 部署较复杂 |
| 多队列 | 简单场景 | 实现简单 | 功能有限 |
最佳实践建议
- 定义清晰的优先级级别:通常3-5级(紧急、高、中、低、后台)
- 防止优先级饥饿:低优先级任务可能永远得不到处理
- 动态调整优先级:根据等待时间自动提升优先级
- 监控告警:监控高优先级任务的处理延迟
- 错误重试机制:失败任务自动提升优先级重试
// 防止优先级饥饿的示例
function getTaskWithAging(PriorityQueue $queue) {
$task = $queue->dequeue();
if ($task) {
// 检查等待时间
$waitTime = time() - $task['created_at'];
// 等待超过5分钟的任务自动提升优先级
if ($waitTime > 300) {
$task['priority'] = max(1, $task['priority'] - 1);
$queue->enqueue($task, $task['priority']);
return getTaskWithAging($queue); // 递归获取其他任务
}
}
return $task;
}
选择哪种方案取决于你的具体需求:任务量、是否需要分布式、持久化要求、预算等因素。