本文目录导读:

我将为您提供几种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方案开始,随着业务发展逐步升级到更可靠的消息队列方案。