PHP实现延迟队列方案

wen PHP项目 1

本文目录导读:

PHP实现延迟队列方案

  1. 方案一:基于Redis的延迟队列
  2. 方案二:基于MySQL的延迟队列
  3. 方案三:基于RabbitMQ插件(推荐)
  4. 方案四:基于Laravel的延迟队列
  5. 方案选择建议
  6. 生产环境建议

我将为您提供几种PHP实现延迟队列的方案,从简单到复杂,您可以根据实际需求选择。

基于Redis的延迟队列

使用Redis ZSet实现(最简单)

<?php
/**
 * 基于Redis ZSet的延迟队列
 */
class RedisDelayQueue
{
    private $redis;
    private $queuePrefix = 'delay_queue:';
    public function __construct($host = '127.0.0.1', $port = 6379)
    {
        $this->redis = new Redis();
        $this->redis->connect($host, $port);
    }
    /**
     * 添加延迟任务
     * @param string $queue 队列名称
     * @param mixed $data 任务数据
     * @param int $delay 延迟秒数
     */
    public function addTask($queue, $data, $delay)
    {
        $score = time() + $delay;
        return $this->redis->zAdd($this->queuePrefix . $queue, $score, json_encode([
            'data' => $data,
            'create_time' => time()
        ]));
    }
    /**
     * 获取到期的任务
     * @param string $queue 队列名称
     */
    public function getExpiredTasks($queue)
    {
        $currentTime = time();
        $tasks = $this->redis->zRangeByScore(
            $this->queuePrefix . $queue,
            0, 
            $currentTime,
            ['limit' => [0, 100]]
        );
        if (!empty($tasks)) {
            // 移除已获取的任务
            $this->redis->zRem($this->queuePrefix . $queue, ...$tasks);
        }
        return array_map(function($task) {
            return json_decode($task, true);
        }, $tasks);
    }
}
// 使用示例
$queue = new RedisDelayQueue();
// 添加一个5秒后执行的任务
$queue->addTask('order_timeout', ['order_id' => 123456], 5);
// 消费者获取到期的任务
$tasks = $queue->getExpiredTasks('order_timeout');
foreach ($tasks as $task) {
    // 处理业务逻辑
    echo "处理订单:" . $task['data']['order_id'] . PHP_EOL;
}

Redis键空间通知实现

<?php
/**
 * 使用Redis键空间通知(Key Space Notification)实现
 */
class RedisNotifyDelayQueue
{
    private $redis;
    private $queuePrefix = 'delay_task:';
    public function __construct($host = '127.0.0.1', $port = 6379)
    {
        $this->redis = new Redis();
        $this->redis->connect($host, $port);
        // 启用键空间通知,需要修改redis.conf
        // notify-keyspace-events "Ex"
    }
    /**
     * 添加延迟任务
     */
    public function addTask($key, $data, $delay)
    {
        $redisKey = $this->queuePrefix . $key;
        // 使用expire实现自动过期
        $this->redis->setex($redisKey, $delay, json_encode($data));
        // 订阅过期事件
        $this->subscribeExpiredEvent($key);
    }
    /**
     * 订阅过期事件(伪代码,需要在单独的进程运行)
     */
    private function subscribeExpiredEvent($key)
    {
        $pubsub = $this->redis->pubSubLoop();
        $pubsub->subscribe('__keyevent@0__:expired');
        foreach ($pubsub as $message) {
            if ($message->kind === 'message' && $message->payload == $this->queuePrefix . $key) {
                // 处理过期任务
                $data = json_decode($this->getOriginalData($key), true);
                $this->processTask($data);
            }
        }
    }
}

基于MySQL的延迟队列

<?php
/**
 * 基于MySQL的延迟队列
 */
class MysqlDelayQueue
{
    private $pdo;
    public function __construct($dsn, $username, $password)
    {
        $this->pdo = new PDO($dsn, $username, $password);
        $this->createTable();
    }
    /**
     * 创建任务表
     */
    private function createTable()
    {
        $sql = "CREATE TABLE IF NOT EXISTS delay_queue (
            id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
            queue_name VARCHAR(100) NOT NULL,
            data TEXT NOT NULL,
            available_time INT UNSIGNED NOT NULL,
            status TINYINT NOT NULL DEFAULT 0,
            created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
            INDEX (queue_name, available_time, status)
        ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;";
        $this->pdo->exec($sql);
    }
    /**
     * 添加延迟任务
     */
    public function addTask($queueName, $data, $delay)
    {
        $sql = "INSERT INTO delay_queue (queue_name, data, available_time) VALUES (?, ?, ?)";
        $stmt = $this->pdo->prepare($sql);
        $stmt->execute([
            $queueName,
            json_encode($data),
            time() + $delay
        ]);
        return $this->pdo->lastInsertId();
    }
    /**
     * 获取到期的任务
     */
    public function getExpiredTasks($queueName, $limit = 100)
    {
        // 使用事务和行锁避免并发
        $this->pdo->beginTransaction();
        $sql = "SELECT * FROM delay_queue 
                WHERE queue_name = ? AND status = 0 AND available_time <= ? 
                ORDER BY available_time ASC 
                LIMIT ? 
                FOR UPDATE";
        $stmt = $this->pdo->prepare($sql);
        $stmt->execute([$queueName, time(), $limit]);
        $tasks = $stmt->fetchAll(PDO::FETCH_ASSOC);
        if (!empty($tasks)) {
            // 标记任务为处理中
            $ids = array_column($tasks, 'id');
            $inClause = implode(',', array_fill(0, count($ids), '?'));
            $updateSql = "UPDATE delay_queue SET status = 1 WHERE id IN ($inClause)";
            $updateStmt = $this->pdo->prepare($updateSql);
            $updateStmt->execute($ids);
        }
        $this->pdo->commit();
        return array_map(function($task) {
            $task['data'] = json_decode($task['data'], true);
            return $task;
        }, $tasks);
    }
    /**
     * 任务完成
     */
    public function completeTask($taskId)
    {
        $sql = "UPDATE delay_queue SET status = 2 WHERE id = ?";
        $stmt = $this->pdo->prepare($sql);
        return $stmt->execute([$taskId]);
    }
}

基于RabbitMQ插件(推荐)

<?php
/**
 * 使用RabbitMQ延迟插件(rabbitmq_delayed_message_exchange)
 */
class RabbitMQDelayQueue
{
    private $connection;
    private $channel;
    public function __construct($host, $port, $user, $pass)
    {
        $this->connection = new AMQPConnection([
            'host' => $host,
            'port' => $port,
            'login' => $user,
            'password' => $pass
        ]);
        $this->connection->connect();
        $this->channel = new AMQPChannel($this->connection);
    }
    /**
     * 创建延迟交换机
     */
    public function createDelayedExchange($exchangeName)
    {
        $exchange = new AMQPExchange($this->channel);
        $exchange->setName($exchangeName);
        $exchange->setType('x-delayed-message');
        $exchange->setArguments([
            'x-delayed-type' => 'direct'
        ]);
        $exchange->declareExchange();
        return $exchange;
    }
    /**
     * 添加延迟任务
     */
    public function addTask($exchangeName, $queueName, $data, $delay)
    {
        $exchange = $this->createDelayedExchange($exchangeName);
        // 创建队列
        $queue = new AMQPQueue($this->channel);
        $queue->setName($queueName);
        $queue->declareQueue();
        $queue->bind($exchangeName, $queueName);
        // 设置消息属性
        $message = new AMQPMessage(json_encode($data), [
            'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
            'headers' => [
                'x-delay' => $delay * 1000 // 毫秒
            ]
        ]);
        return $exchange->publish($message, $queueName);
    }
    /**
     * 消费任务(消费者)
     */
    public function consume($queueName, $callback)
    {
        $queue = new AMQPQueue($this->channel);
        $queue->setName($queueName);
        $queue->declareQueue();
        while (true) {
            $message = $queue->get();
            if ($message) {
                $data = json_decode($message->getBody(), true);
                $callback($data);
                $queue->ack($message->getDeliveryTag());
            } else {
                usleep(100000); // 100ms
            }
        }
    }
}

基于Laravel的延迟队列

<?php
/**
 * 使用Laravel框架的延迟队列(最易用)
 */
class LaravelDelayQueue
{
    public function addTask()
    {
        // 方法1:使用delay方法
        dispatch(new ProcessOrderJob($order))
            ->delay(now()->addMinutes(10));
        // 方法2:指定等待截止时间
        dispatch(new SendEmailJob($data))
            ->until(now()->addHour());
        // 方法3:使用队列
        $job = new ProcessOrderJob($order);
        $job->delay(now()->addSeconds(30));
        $this->dispatch($job);
        // 方法4:发送到特定队列
        $job = new SendNotificationJob($notification);
        $job->onQueue('notifications')->delay(now()->addMinutes(5));
        $this->dispatch($job);
    }
    /**
     * 实现队列Worker
     */
    public function queueWorker()
    {
        // 命令行执行
        // php artisan queue:work --queue=notifications --tries=3
        // 在队列中处理延迟任务
        // app/Jobs/ProcessOrderJob.php
        /*
        class ProcessOrderJob implements ShouldQueue
        {
            use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
            protected $order;
            public $tries = 3;
            public $timeout = 120;
            public function __construct(Order $order)
            {
                $this->order = $order;
            }
            public function handle()
            {
                // 处理订单超时逻辑
                $this->order->timeout();
            }
            public function failed(Exception $exception)
            {
                // 失败处理
            }
        }
        */
    }
}

方案选择建议

方案 优点 缺点 适用场景
Redis ZSet 简单易实现,高性能 需要持续轮询 中小型系统,延迟精度要求不高
Redis 通知 实时性高 需要Redis配置支持 需要精确触发,且Redis版本支持
MySQL 持久化,易管理 性能受限,需要轮询 数据重要,对实时性要求不高
RabbitMQ插件 高可靠,功能强大 需要额外部署服务 大型系统,高可靠性要求
Laravel Queue 开发效率高,生态完善 依赖框架 基于Laravel的项目

生产环境建议

<?php
/**
 * 生产环境多方案组合使用
 */
class ProductionDelayQueue
{
    private $redisQueue;
    private $mysqlQueue;
    public function __construct()
    {
        // 配置多个队列服务
        $this->redisQueue = new RedisDelayQueue();
        $this->mysqlQueue = new MysqlDelayQueue();
    }
    /**
     * 添加延迟任务(带重试机制)
     */
    public function addTaskWithRetry($queueName, $data, $delay, $retryCount = 3)
    {
        $result = false;
        $attempt = 0;
        while ($attempt < $retryCount && !$result) {
            try {
                // 优先使用Redis
                $result = $this->redisQueue->addTask($queueName, $data, $delay);
                // 如果Redis失败,使用MySQL
                if (!$result) {
                    $result = $this->mysqlQueue->addTask($queueName, $data, $delay);
                }
            } catch (Exception $e) {
                // 记录错误日志
                error_log($e->getMessage());
            }
            $attempt++;
        }
        return $result;
    }
}

选择哪种方案取决于您的具体需求:系统规模、可靠性要求、开发成本和运维复杂度,建议从简单的Redis方案开始,随着业务发展逐步升级到更可靠的消息队列方案。

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