PHP项目转码队列如何优先级排序任务

wen PHP项目 30

本文目录导读:

PHP项目转码队列如何优先级排序任务

  1. 数据库队列方案(推荐)
  2. Redis 有序集合方案(高性能)
  3. 消息队列扩展方案(RabbitMQ)
  4. 混合方案(推荐生产环境)
  5. 性能优化建议

在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);
    }
}

性能优化建议

  1. 分页处理:批量拉取多个高优先级任务
  2. 超时机制:对长时间处理中的任务进行重试
  3. 监控告警:队列长度、处理时间监控
  4. 并发控制:根据服务器资源调整并发数
// 批量拉取任务示例
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 + 持久化
  • 生产环境:推荐混合方案

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