本文目录导读:

实现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项目 |
最佳实践建议
- 数据持久化:重要任务务必存储在数据库
- 幂等性设计:确保任务重复执行不影响业务
- 监控告警:监控队列积压情况
- 降级方案:主服务故障时备用方案
- 清理机制:定期清理已完成的任务
选择方案时,建议根据项目规模、并发量、可靠性要求综合考虑,对于大多数PHP项目,推荐使用Redis ZSet + 数据库兜底的方案。