PHP项目延时队列如何实现定时触发任务

wen PHP项目 27

本文目录导读:

PHP项目延时队列如何实现定时触发任务

  1. 基于Redis的延时队列
  2. 使用消息队列中间件
  3. 使用PHP扩展或框架
  4. 自建MySQL延时轮询
  5. 完整示例:订单超时取消
  6. 定时触发脚本(Crontab)
  7. 方案对比
  8. 最佳实践建议

实现PHP项目的延时队列定时触发任务,常见有几种方案,我来详细说明。

基于Redis的延时队列

使用Redis的有序集合(ZSet)

<?php
class RedisDelayQueue
{
    private $redis;
    private $queueKey = 'delay_queue';
    public function __construct($redis)
    {
        $this->redis = $redis;
    }
    /**
     * 添加延时任务
     * @param string $taskId 任务ID
     * @param array $data 任务数据
     * @param int $delay 延迟秒数
     */
    public function addTask($taskId, $data, $delay)
    {
        $score = time() + $delay; // 执行时间戳
        $this->redis->zAdd($this->queueKey, $score, json_encode([
            'task_id' => $taskId,
            'data' => $data,
            'create_time' => time()
        ]));
    }
    /**
     * 获取到期任务
     */
    public function getExpiredTasks()
    {
        $now = time();
        // 获取当前时间之前的所有任务
        $tasks = $this->redis->zRangeByScore($this->queueKey, 0, $now);
        if (!empty($tasks)) {
            // 移除这些任务
            $this->redis->zRemRangeByScore($this->queueKey, 0, $now);
        }
        return $tasks;
    }
}
// 使用示例
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$queue = new RedisDelayQueue($redis);
// 添加延时任务,30分钟后执行
$queue->addTask('order_123', [
    'type' => 'order_timeout',
    'order_id' => '12345'
], 1800);

基于Redis Keyspace Notifications

<?php
// 设置Redis配置支持过期通知
// config set notify-keyspace-events Ex
class RedisKeyExpireQueue
{
    private $redis;
    private $prefix = 'delay_task:';
    public function __construct($redis)
    {
        $this->redis = $redis;
    }
    /**
     * 添加延时任务
     */
    public function addTask($taskType, $taskData, $delaySeconds)
    {
        $taskId = uniqid();
        $key = $this->prefix . $taskId;
        // 存储任务数据
        $this->redis->set($key, json_encode([
            'type' => $taskType,
            'data' => $taskData
        ]));
        // 设置过期时间
        $this->redis->expire($key, $delaySeconds);
        return $taskId;
    }
    /**
     * 监听过期事件(需要后台脚本运行)
     */
    public function listenExpiredEvents()
    {
        $this->redis->psubscribe(['__keyevent@0__:expired'], function($redis, $pattern, $channel, $message) {
            echo "Task expired: {$message}\n";
            // 获取并处理过期任务
            $this->handleExpiredTask($message);
        });
    }
    private function handleExpiredTask($key)
    {
        // 处理过期的任务
        // 注意:此时key已过期,数据无法获取
        // 建议在实际业务中另存任务数据
    }
}

使用消息队列中间件

RabbitMQ 延时队列

<?php
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
use PhpAmqpLib\Wire\AMQPTable;
class RabbitMQDelayQueue
{
    private $connection;
    private $channel;
    public function __construct($host, $port, $user, $password)
    {
        $this->connection = new AMQPStreamConnection($host, $port, $user, $password);
        $this->channel = $this->connection->channel();
    }
    /**
     * 设置延时队列
     */
    public function setupDelayQueue()
    {
        // 交换机
        $exchange = 'delay_exchange';
        $this->channel->exchange_declare($exchange, 'direct', false, true, false);
        // 正常队列
        $queue = 'delay_queue';
        $this->channel->queue_declare($queue, false, true, false, false);
        // 死信队列设置
        $args = new AMQPTable([
            'x-dead-letter-exchange' => 'delay_exchange',
            'x-dead-letter-routing-key' => 'process_queue',
            'x-message-ttl' => 60000 // 默认延时60秒
        ]);
        $delayQueue = 'delay_queue_ttl';
        $this->channel->queue_declare($delayQueue, false, true, false, false, false, false, $args);
        // 绑定
        $this->channel->queue_bind($delayQueue, $exchange, 'delay');
        return [$exchange, $queue];
    }
    /**
     * 发送延时消息
     */
    public function sendDelayMessage($data, $delayMs)
    {
        $msg = new AMQPMessage(json_encode($data), [
            'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT
        ]);
        // 设置消息级别的TTL
        $headers = new AMQPTable([
            'x-delay' => $delayMs
        ]);
        $msg->set('application_headers', $headers);
        $this->channel->basic_publish($msg, 'delay_exchange', 'delay');
    }
    public function __destruct()
    {
        $this->channel->close();
        $this->connection->close();
    }
}

Beanbalkd 简单轻量级

<?php
class BeanstalkdQueue
{
    private $pheanstalk;
    public function __construct($host = '127.0.0.1', $port = 11300)
    {
        $this->pheanstalk = \Pheanstalk\Pheanstalk::create($host, $port);
    }
    /**
     * 添加延时任务
     */
    public function addDelayedTask($tube, $data, $delay)
    {
        $this->pheanstalk
            ->useTube($tube)
            ->put(json_encode($data), 
                \Pheanstalk\Pheanstalk::DEFAULT_PRIORITY, // 优先级
                $delay, // 延迟时间(秒)
                \Pheanstalk\Pheanstalk::DEFAULT_TTR // 处理时间
            );
    }
    /**
     * 消费任务(守护进程运行)
     */
    public function consume($tube, callable $callback)
    {
        while (true) {
            $job = $this->pheanstalk
                ->watch($tube)
                ->ignore('default')
                ->reserve();
            if ($job) {
                try {
                    $data = json_decode($job->getData(), true);
                    $callback($data);
                    $this->pheanstalk->delete($job);
                } catch (\Exception $e) {
                    // 处理失败,可以 bury 或 release
                    $this->pheanstalk->bury($job);
                }
            }
            sleep(1);
        }
    }
}

使用PHP扩展或框架

Laravel 任务调度

<?php
// 定义任务类
namespace App\Jobs;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;
class ProcessOrder implements ShouldQueue
{
    use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
    protected $orderId;
    public function __construct($orderId)
    {
        $this->orderId = $orderId;
    }
    public function handle()
    {
        // 处理超时订单
        Order::where('id', $this->orderId)
            ->where('status', 'pending')
            ->update(['status' => 'cancelled']);
    }
}
// 添加延时任务(延迟30分钟)
ProcessOrder::dispatch($orderId)->delay(now()->addMinutes(30));

ThinkPHP 队列

<?php
namespace app\job;
use think\queue\Job;
class DelayTask
{
    public function fire(Job $job, $data)
    {
        // 处理任务
        $orderId = $data['order_id'];
        // 业务逻辑...
        // 任务执行成功
        $job->delete();
    }
}
// 添加延时任务
Queue::later(3600, 'app\job\DelayTask', ['order_id' => 123], 'queue_name');

自建MySQL延时轮询

<?php
class MySQLDelayQueue
{
    private $db;
    public function __construct($db)
    {
        $this->db = $db;
    }
    /**
     * 添加延时任务
     */
    public function addTask($taskType, $taskData, $delaySeconds)
    {
        $executeTime = date('Y-m-d H:i:s', time() + $delaySeconds);
        $stmt = $this->db->prepare(
            "INSERT INTO delay_queue (task_type, task_data, execute_time, status, create_time) 
             VALUES (?, ?, ?, 'pending', NOW())"
        );
        $stmt->execute([$taskType, json_encode($taskData), $executeTime]);
        return $this->db->lastInsertId();
    }
    /**
     * 获取并执行到期任务
     */
    public function processExpiredTasks()
    {
        // 开启事务
        $this->db->beginTransaction();
        try {
            // 获取到期的待处理任务
            $stmt = $this->db->prepare(
                "SELECT * FROM delay_queue 
                 WHERE execute_time <= NOW() 
                 AND status = 'pending' 
                 LIMIT 10 FOR UPDATE"
            );
            $stmt->execute();
            $tasks = $stmt->fetchAll();
            foreach ($tasks as $task) {
                // 标记为处理中
                $updateStmt = $this->db->prepare(
                    "UPDATE delay_queue SET status = 'processing' WHERE id = ?"
                );
                $updateStmt->execute([$task['id']]);
                // 执行任务
                $result = $this->handleTask($task);
                // 更新状态
                $finalStatus = $result ? 'completed' : 'failed';
                $updateStmt = $this->db->prepare(
                    "UPDATE delay_queue SET status = ?, complete_time = NOW() WHERE id = ?"
                );
                $updateStmt->execute([$finalStatus, $task['id']]);
            }
            $this->db->commit();
        } catch (\Exception $e) {
            $this->db->rollBack();
            throw $e;
        }
    }
    private function handleTask($task)
    {
        $taskData = json_decode($task['task_data'], true);
        // 根据任务类型处理
        switch ($task['task_type']) {
            case 'order_timeout':
                // 取消订单
                return $this->cancelOrder($taskData['order_id']);
            default:
                return false;
        }
    }
}

完整示例:订单超时取消

<?php
/**
 * 综合使用Redis + 定时脚本实现订单超时取消
 */
class OrderTimeoutSystem
{
    private $redis;
    private $db;
    private $queuePrefix = 'order_timeout:';
    public function __construct($redis, $db)
    {
        $this->redis = $redis;
        $this->db = $db;
    }
    /**
     * 创建订单并设置超时
     */
    public function createOrder($userId, $productId, $price)
    {
        // 1. 创建订单
        $orderId = $this->createOrderRecord($userId, $productId, $price);
        // 2. 设置延时任务(30分钟超时)
        $this->setTimeoutTask($orderId, 1800);
        return $orderId;
    }
    /**
     * 设置超时任务
     */
    private function setTimeoutTask($orderId, $timeout)
    {
        $key = $this->queuePrefix . $orderId;
        // 使用Redis SETNX防止重复设置
        if ($this->redis->setnx($key, '1')) {
            $this->redis->expire($key, $timeout);
        }
        // 同时记录到有序集合,便于管理和监控
        $this->redis->zAdd('order_timeout_zset', time() + $timeout, $orderId);
    }
    /**
     * 处理超时订单(由定时任务每分钟调用)
     */
    public function processTimeoutOrders()
    {
        // 1. 检查Redis过期通知(主动方式)
        $expiredOrders = $this->redis->keys($this->queuePrefix . '*');
        foreach ($expiredOrders as $key) {
            $orderId = str_replace($this->queuePrefix, '', $key);
            $this->cancelTimeoutOrder($orderId);
        }
        // 2. 数据库兜底检查(防止Redis数据丢失)
        $this->dbBackupCheck();
    }
    /**
     * 取消超时订单
     */
    private function cancelTimeoutOrder($orderId)
    {
        // 加锁防止并发
        $lockKey = "lock:order:{$orderId}";
        if ($this->redis->setnx($lockKey, time())) {
            $this->redis->expire($lockKey, 10);
            try {
                // 检查订单状态
                $order = $this->db->query(
                    "SELECT status FROM orders WHERE id = ?",
                    [$orderId]
                );
                if ($order && $order['status'] === 'pending') {
                    // 取消订单
                    $this->db->execute(
                        "UPDATE orders SET status = 'cancelled', cancel_time = NOW(), 
                         cancel_reason = 'timeout' WHERE id = ?",
                        [$orderId]
                    );
                    // 释放库存等业务逻辑
                    $this->releaseStock($orderId);
                    // 清理Redis中的任务
                    $this->redis->del($this->queuePrefix . $orderId);
                    $this->redis->zRem('order_timeout_zset', $orderId);
                }
            } finally {
                $this->redis->del($lockKey);
            }
        }
    }
    /**
     * 数据库兜底检查
     */
    private function dbBackupCheck()
    {
        // 查询超时未处理的订单
        $timeoutOrders = $this->db->query(
            "SELECT id FROM orders 
             WHERE status = 'pending' 
             AND create_time < DATE_SUB(NOW(), INTERVAL 30 MINUTE)
             AND id NOT IN (
                 SELECT order_id FROM order_timeout_log WHERE processed = 1
             )
             LIMIT 100"
        );
        foreach ($timeoutOrders as $order) {
            $this->cancelTimeoutOrder($order['id']);
        }
    }
}

定时触发脚本(Crontab)

# 每分钟执行一次延时队列处理
* * * * * /usr/bin/php /path/to/process_delay_queue.php >> /var/log/delay_queue.log 2>&1
# 每5分钟执行一次
*/5 * * * * /usr/bin/php /path/to/process_delay_queue.php

方案对比

方案 优点 缺点 适用场景
Redis ZSet 速度快、实现简单 数据可能丢失 中小型项目
RabbitMQ 可靠、支持ACK 部署复杂 大型系统
Beanstalkd 轻量、易用 功能有限 中等规模
MySQL轮询 数据安全 性能差 简单场景
Laravel队列 功能完善 依赖框架 Laravel项目

最佳实践建议

  1. 数据持久化:重要任务务必存储在数据库
  2. 幂等性设计:确保任务重复执行不影响业务
  3. 监控告警:监控队列积压情况
  4. 降级方案:主服务故障时备用方案
  5. 清理机制:定期清理已完成的任务

选择方案时,建议根据项目规模、并发量、可靠性要求综合考虑,对于大多数PHP项目,推荐使用Redis ZSet + 数据库兜底的方案。

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