本文目录导读:

在PHP项目中实现转码队列的优先级排序,可以采用以下几种方案,从简单到复杂:
数据库队列方案(推荐)
表结构设计
CREATE TABLE transcode_queue (
id INT PRIMARY KEY AUTO_INCREMENT,
file_path VARCHAR(500) NOT NULL,
target_format VARCHAR(20) NOT NULL,
priority INT NOT NULL DEFAULT 0, -- 优先级数字越小越优先
status ENUM('pending', 'processing', 'completed', 'failed') DEFAULT 'pending',
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
retry_count INT DEFAULT 0,
timeout_at TIMESTAMP,
INDEX idx_priority_status (priority, status),
INDEX idx_status (status)
);
任务入队
class TranscodeQueue {
private PDO $pdo;
public function enqueue(string $filePath, string $targetFormat, int $priority = 0): int {
$stmt = $this->pdo->prepare(
"INSERT INTO transcode_queue (file_path, target_format, priority)
VALUES (:file_path, :target_format, :priority)"
);
$stmt->execute([
':file_path' => $filePath,
':target_format' => $targetFormat,
':priority' => $priority
]);
return $this->pdo->lastInsertId();
}
public function getNextTask(): ?array {
$this->pdo->beginTransaction();
try {
// 获取优先级最高的待处理任务
$stmt = $this->pdo->prepare(
"SELECT * FROM transcode_queue
WHERE status = 'pending'
ORDER BY priority ASC, created_at ASC
LIMIT 1 FOR UPDATE SKIP LOCKED"
);
$stmt->execute();
$task = $stmt->fetch(PDO::FETCH_ASSOC);
if ($task) {
// 标记为处理中
$update = $this->pdo->prepare(
"UPDATE transcode_queue SET status = 'processing' WHERE id = :id"
);
$update->execute([':id' => $task['id']]);
}
$this->pdo->commit();
return $task;
} catch (Exception $e) {
$this->pdo->rollBack();
return null;
}
}
}
Redis 有序集合方案(高性能)
class RedisPriorityQueue {
private Redis $redis;
private string $queueKey = 'transcode:queue';
private string $processingKey = 'transcode:processing';
public function enqueue(string $taskId, array $taskData, int $priority = 0): void {
// 优先级数字越小越优先,所以用负数
$score = -$priority . '.' . microtime(true);
$this->redis->zAdd($this->queueKey, $score, json_encode([
'id' => $taskId,
'data' => $taskData,
'priority' => $priority
]));
}
public function getNextTask(): ?array {
// 获取优先级最高的任务
$tasks = $this->redis->zPopMin($this->queueKey, 1);
if (empty($tasks)) {
return null;
}
$task = json_decode($tasks[0], true);
// 移动到处理中队列
$this->redis->hSet($this->processingKey, $task['id'], json_encode($task));
return $task;
}
public function completeTask(string $taskId): void {
$this->redis->hDel($this->processingKey, $taskId);
}
// 获取队列统计信息
public function getStats(): array {
return [
'queue_length' => $this->redis->zCard($this->queueKey),
'processing_count' => $this->redis->hLen($this->processingKey)
];
}
}
消息队列扩展方案(RabbitMQ)
安装 AMQP 扩展
composer require php-amqplib/php-amqplib
实现代码
class RabbitMQPriorityQueue {
private AMQPConnection $connection;
private AMQPChannel $channel;
public function enqueue(string $taskId, array $data, int $priority = 0): void {
$msg = new AMQPMessage(json_encode([
'id' => $taskId,
'data' => $data
]), [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
'priority' => $priority
]);
// 创建优先级队列
$channel->queue_declare(
'transcode_queue',
false,
true, // durable
false,
false,
false,
new AMQPTable([
'x-max-priority' => 10 // 优先级别 0-10
])
);
$channel->basic_publish($msg, '', 'transcode_queue');
}
public function worker(): void {
$channel->basic_qos(null, 1, null);
$channel->basic_consume(
'transcode_queue',
'',
false,
false,
false,
false,
function($msg) {
$task = json_decode($msg->body, true);
// 处理任务
$this->processTask($task);
$msg->ack();
}
);
while ($channel->is_consuming()) {
$channel->wait();
}
}
}
混合方案(推荐生产环境)
结合数据库持久化和Redis优先级排序:
class HybridPriorityQueue {
private PDO $pdo;
private Redis $redis;
public function enqueue($task): void {
// 入库持久化
$taskId = $this->saveToDatabase($task);
// 更新Redis排序缓存
$this->redis->zAdd('priority_index',
-$task['priority'],
$taskId
);
}
public function getNextTask(): ?array {
// 从Redis获取最高优先级任务ID
$taskId = $this->redis->zPopMin('priority_index', 1);
if (!$taskId) {
// Redis没有则从数据库获取
return $this->getTaskFromDB();
}
// 从数据库获取完整任务信息
return $this->getTaskById($taskId);
}
}
性能优化建议
- 分页处理:批量拉取多个高优先级任务
- 超时机制:对长时间处理中的任务进行重试
- 监控告警:队列长度、处理时间监控
- 并发控制:根据服务器资源调整并发数
// 批量拉取任务示例
public function getBatchTasks(int $limit = 5): array {
$stmt = $this->pdo->query(
"SELECT * FROM transcode_queue
WHERE status = 'pending'
ORDER BY priority ASC, created_at ASC
LIMIT $limit"
);
$tasks = $stmt->fetchAll();
// 批量更新状态
$ids = array_column($tasks, 'id');
$this->pdo->exec(
"UPDATE transcode_queue
SET status = 'processing',
started_at = NOW()
WHERE id IN (" . implode(',', $ids) . ")"
);
return $tasks;
}
选择合适的方案取决于你的具体需求:
- 简单项目:数据库队列即可
- 高并发:Redis方案
- 企业级:RabbitMQ + 持久化
- 生产环境:推荐混合方案