在高并发 Web 应用中,同步处理耗时任务会严重拖慢响应速度。发送邮件、生成报表、处理视频转码——这些操作动辄耗时数秒甚至数分钟,如果在用户请求中同步执行,必然导致请求超时和用户体验下降。消息队列是解决这一问题的核心架构模式,它将任务从请求主线程中剥离,交给后台 Worker 异步消费,从而实现请求的快速返回和任务的可靠处理。
本文将从零开始,系统讲解 PHP 项目中消息队列的架构设计与实战实现。我们将从最简单的 Redis 队列入手,逐步演进到 RabbitMQ 的生产级方案,覆盖队列选型、Worker 管理、失败重试、监控告警等关键环节,并提供可直接使用的代码示例。
一、消息队列的核心概念与选型
消息队列(Message Queue)的本质是一个生产者-消费者模型:生产者将任务封装为消息投递到队列,消费者从队列中取出消息并执行。这种解耦带来了三大核心价值:
- 异步化:耗时操作不阻塞用户请求,响应时间从分钟级降到毫秒级
- 削峰填谷:突发流量时任务排队等待,Worker 按自身节奏消费,保护后端服务
- 解耦:生产者不需要知道消费者是谁、在哪里、有多少个,只需把消息扔进队列
1.1 Redis vs RabbitMQ 选型对比
| 特性 | Redis 队列 | RabbitMQ |
|---|---|---|
| 部署复杂度 | 低(已有 Redis 即可) | 中(需独立部署 Erlang 服务) |
| 消息持久化 | 可选(AOF/RDB) | 原生支持 |
| 消息确认 | 需自行实现 | 原生 ACK 机制 |
| 路由能力 | 无(纯列表模型) | Exchange + Binding 灵活路由 |
| 死信队列 | 需自行实现 | 原生 DLX 支持 |
| 吞吐量 | 极高(10万+/秒) | 高(万级/秒) |
| 适用场景 | 轻量异步、任务简单 | 复杂路由、可靠投递、多消费者 |
经验法则:项目初期用 Redis 快速上线,当需要消息确认、复杂路由、死信处理等高级特性时再迁移到 RabbitMQ。两者并非互斥——很多生产系统同时使用 Redis 做缓存和 RabbitMQ 做消息队列。
二、Redis 队列实战:轻量级异步方案
Redis 的 List 数据结构天然适合做队列——
1 | LPUSH |
入队、
1 | BRPOP |
出队,两条命令就能实现一个基本队列。PHP 中通过
1 | phpredis |
扩展可以高效操作。
2.1 队列封装
我们先将 Redis 队列操作封装为一个可复用的类,支持多队列、延迟任务和优先级控制:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125 <?php
// RedisQueue.php
class RedisQueue
{
private Redis $redis;
private string $prefix;
public function __construct(Redis $redis, string $prefix = 'queue:')
{
$this->redis = $redis;
$this->prefix = $prefix;
}
/**
* 投递即时任务
*/
public function push(string $queue, array $payload): void
{
$message = json_encode([
'id' => uniqid('msg_', true),
'payload' => $payload,
'attempts' => 0,
'created' => time(),
]);
$this->redis->lPush($this->prefix . $queue, $message);
}
/**
* 投递延迟任务(使用 Sorted Set + 时间戳分数)
*/
public function later(string $queue, array $payload, int $delay): void
{
$message = json_encode([
'id' => uniqid('msg_', true),
'payload' => $payload,
'attempts' => 0,
'created' => time(),
'run_at' => time() + $delay,
]);
$this->redis->zAdd(
$this->prefix . $queue . ':delayed',
time() + $delay,
$message
);
}
/**
* 消费任务(阻塞式,超时返回 null)
*/
public function pop(string $queue, int $timeout = 30): ?array
{
// 先迁移到期的延迟任务
$this->migrateDelayed($queue);
$result = $this->redis->brPop(
[$this->prefix . $queue],
$timeout
);
if (!$result) {
return null;
}
$message = json_decode($result[1], true);
// 移入处理中列表,实现至少一次消费
$this->redis->zAdd(
$this->prefix . $queue . ':reserved',
time() + 300,
$result[1]
);
return $message;
}
/**
* 确认任务完成
*/
public function ack(string $queue, string $rawMessage): void
{
$this->redis->zRem(
$this->prefix . $queue . ':reserved',
$rawMessage
);
}
/**
* 任务失败,重新入队或移入死信
*/
public function fail(string $queue, string $rawMessage, int $maxRetries = 3): void
{
$this->redis->zRem(
$this->prefix . $queue . ':reserved',
$rawMessage
);
$message = json_decode($rawMessage, true);
$message['attempts']++;
if ($message['attempts'] >= $maxRetries) {
$this->redis->lPush(
$this->prefix . $queue . ':failed',
json_encode($message)
);
} else {
$delay = pow(2, $message['attempts']) * 10;
$this->redis->zAdd(
$this->prefix . $queue . ':delayed',
time() + $delay,
json_encode($message)
);
}
}
private function migrateDelayed(string $queue): void
{
$now = time();
$key = $this->prefix . $queue . ':delayed';
$messages = $this->redis->zRangeByScore($key, '-inf', (string)$now);
foreach ($messages as $message) {
$this->redis->zRem($key, $message);
$this->redis->lPush($this->prefix . $queue, $message);
}
}
}
这个封装实现了消息 ID、重试计数、延迟任务迁移、处理超时回收和死信队列——虽然不如 RabbitMQ 的原生机制优雅,但在 Redis 生态内足够可靠。
2.2 Worker 进程实现
Worker 是常驻进程,循环从队列取任务执行。我们需要处理优雅退出、内存泄漏和超时控制:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53 <?php
// worker.php
require __DIR__ . '/RedisQueue.php';
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$queue = new RedisQueue($redis);
$shutdown = false;
pcntl_async_signals(true);
pcntl_signal(SIGTERM, fn() => $shutdown = true);
pcntl_signal(SIGINT, fn() => $shutdown = true);
$memoryLimit = 128 * 1024 * 1024; // 128MB
$jobsProcessed = 0;
echo "Worker started, listening on queue: emails\n";
while (!$shutdown) {
if (memory_get_usage(true) > $memoryLimit) {
echo "Memory limit reached, exiting for restart\n";
break;
}
$job = $queue->pop('emails', 30);
if (!$job) {
continue;
}
$rawMessage = json_encode($job);
try {
$result = dispatch($job['payload']);
$queue->ack('emails', $rawMessage);
$jobsProcessed++;
echo "[OK] Job {$job['id']} processed\n";
} catch (Throwable $e) {
$queue->fail('emails', $rawMessage, 3);
echo "[FAIL] Job {$job['id']}: {$e->getMessage()}\n";
}
}
echo "Worker shutdown, processed {$jobsProcessed} jobs\n";
function dispatch(array $payload): mixed
{
return match ($payload['type']) {
'send_email' => sendEmail($payload),
'generate_report' => generateReport($payload),
default => throw new RuntimeException("Unknown type"),
};
}
2.3 Supervisor 管理 Worker 进程
Worker 进程必须由进程管理器守护,确保崩溃后自动重启。Supervisor 是 PHP 生态的标准选择:
1
2
3
4
5
6
7
8
9
10
11
12
13
14 ; /etc/supervisor/conf.d/worker-emails.conf
[program:worker-emails]
command=php /var/www/worker.php
process_name=%(program_name)s_%(process_num)02d
numprocs=4
autostart=true
autorestart=true
startsecs=3
stopwaitsecs=10
user=www-data
redirect_stderr=true
stdout_logfile=/var/log/worker-emails.log
stdout_logfile_maxbytes=50MB
stdout_logfile_backups=5
关键参数说明:
-
1numprocs=4
:启动 4 个 Worker 并行消费,根据任务 CPU/IO 特性调整
-
1autorestart=true
:进程退出后立即重启,包括正常退出(解决内存泄漏)
-
1stopwaitsecs=10
:优雅停止等待时间,Worker 在此期间完成当前任务后退出
三、RabbitMQ 队列实战:生产级可靠方案
当项目需要消息持久化、复杂路由、优先级队列和原生 ACK 机制时,RabbitMQ 是更专业的选择。它基于 AMQP 协议,提供了 Exchange(交换机)、Binding(绑定)、DLX(死信交换)等完整的消息模型。
3.1 连接与通道管理
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27 <?php
// RabbitMQConnection.php
use PhpAmqpLib\Connection\AMQPStreamConnection;
class RabbitMQConnection
{
private static ?AMQPStreamConnection $connection = null;
public static function getConnection(): AMQPStreamConnection
{
if (self::$connection === null || self::$connection->isClosed()) {
self::$connection = new AMQPStreamConnection(
host: getenv('RABBITMQ_HOST') ?: '127.0.0.1',
port: (int)(getenv('RABBITMQ_PORT') ?: 5672),
user: getenv('RABBITMQ_USER') ?: 'guest',
password: getenv('RABBITMQ_PASS') ?: 'guest',
vhost: getenv('RABBITMQ_VHOST') ?: '/'
);
}
return self::$connection;
}
public static function createChannel()
{
return self::getConnection()->channel();
}
}
注意:AMQP 连接是重量级资源,应该复用;Channel 是轻量级的,每个 Worker 使用独立 Channel。
3.2 交换机与队列声明
RabbitMQ 的路由模型通过 Exchange 实现。我们以
1 | direct |
类型为例,实现按路由键分发任务:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41 <?php
// QueueDeclaration.php
use PhpAmqpLib\Message\AMQPMessage;
use PhpAmqpLib\Wire\AMQPTable;
class QueueDeclaration
{
public static function declareTasks($channel): void
{
// 声明主交换机
$channel->exchange_declare(
'tasks.exchange',
'direct',
false, // passive
true, // durable
false // auto_delete
);
// 声明死信交换机
$channel->exchange_declare(
'tasks.dlx', 'direct', false, true, false
);
// 主队列:绑定死信
$args = new AMQPTable([
'x-dead-letter-exchange' => 'tasks.dlx',
'x-dead-letter-routing-key' => 'failed',
'x-message-ttl' => 86400000,
'x-queue-type' => 'quorum',
]);
$channel->queue_declare(
'tasks.emails',
false, true, false, false, false, $args
);
$channel->queue_bind('tasks.emails', 'tasks.exchange', 'email');
// 死信队列
$channel->queue_declare('tasks.failed', false, true, false, false);
$channel->queue_bind('tasks.failed', 'tasks.dlx', 'failed');
}
}
3.3 生产者:可靠投递
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42 <?php
// TaskProducer.php
use PhpAmqpLib\Message\AMQPMessage;
class TaskProducer
{
private $channel;
public function __construct($channel)
{
$this->channel = $channel;
// 开启发布确认模式
$this->channel->confirm_select();
}
public function publish(
string $routingKey,
array $payload,
int $priority = 0,
int $expiration = 3600000
): void {
$body = json_encode([
'id' => uuid_create(),
'payload' => $payload,
'timestamp' => microtime(true),
]);
$message = new AMQPMessage($body, [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
'priority' => $priority,
'expiration' => (string)$expiration,
'content_type' => 'application/json',
]);
$this->channel->basic_publish(
$message, 'tasks.exchange', $routingKey
);
// 等待确认
$this->channel->wait_for_pending_acks(5);
}
}
1 | confirm_select |
模式确保每条消息都被 Broker 确认写入后才返回。这是生产环境的必选项——没有它,消息可能因为连接断开而丢失,生产者还浑然不知。
3.4 消费者:手动确认与并发控制
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50 <?php
// TaskConsumer.php
use PhpAmqpLib\Message\AMQPMessage;
class TaskConsumer
{
private $channel;
public function __construct($channel)
{
$this->channel = $channel;
}
public function consume(string $queue, callable $handler): void
{
// QoS prefetch = 1,公平调度
$this->channel->basic_qos(null, 1, null);
$callback = function (AMQPMessage $msg) use ($handler) {
$data = json_decode($msg->getBody(), true);
try {
$result = $handler($data['payload']);
$msg->ack();
echo "[ACK] {$data['id']} done\n";
} catch (Throwable $e) {
$headers = $msg->get('application_headers');
$retryCount = $headers
? ($headers->getNativeData()['x-retry-count'] ?? 0)
: 0;
if ($retryCount < 3) {
$msg->nack(false, true);
echo "[RETRY] {$data['id']} attempt " . ($retryCount + 1) . "\n";
} else {
$msg->nack(false, false);
echo "[DLX] {$data['id']} exhausted retries\n";
}
}
};
$this->channel->basic_consume(
$queue, '', false, false, false, false, $callback
);
while ($this->channel->is_consuming()) {
$this->channel->wait(null, false, 30);
}
}
}
关键设计点:
- QoS prefetch = 1:防止某个 Worker 堆积过多未处理消息,实现公平调度
- 手动 ACK:处理成功才确认,异常时
1nack
重新入队或进入死信
- 重试上限:3 次重试后进入死信队列,避免毒药消息无限循环
四、Laravel 队列集成:框架级方案
Laravel 的队列系统是 PHP 生态中最成熟的队列抽象层,开箱支持 Redis、RabbitMQ(通过扩展包)、Database、Amazon SQS 等多种驱动,且内置了重试、延迟、失败处理等机制。
4.1 Job 类定义
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39 <?php
// app/Jobs/SendWelcomeEmail.php
namespace App\Jobs;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;
use App\Models\User;
use App\Services\MailService;
class SendWelcomeEmail implements ShouldQueue
{
use Queueable, InteractsWithQueue, SerializesModels;
public int $tries = 3;
public int $backoff = 60;
public int $timeout = 120;
public bool $deleteWhenMissingModels = true;
public function __construct(
public User $user
) {}
public function handle(MailService $mail): void
{
$mail->send($this->user->email, 'welcome', [
'name' => $this->user->name,
]);
}
public function failed(\Throwable $e): void
{
logger()->error("Welcome email failed", [
'user_id' => $this->user->id,
'error' => $e->getMessage(),
]);
}
}
4.2 分发与 Horizon 监控
1
2
3
4
5
6
7
8
9
10
11
12
13
14 // 控制器中分发任务
use App\Jobs\SendWelcomeEmail;
public function register(Request $request)
{
$user = User::create($request->validated());
SendWelcomeEmail::dispatch($user);
// 指定队列和延迟
SendWelcomeEmail::dispatch($user)
->onQueue('emails')
->delay(now()->addMinutes(5));
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14 // config/horizon.php
'environments' => [
'production' => [
'supervisor-default' => [
'connection' => 'redis',
'queue' => ['default', 'emails', 'reports'],
'balance' => 'auto',
'min_processes' => 2,
'max_processes' => 10,
'balanceMaxShift' => 5,
'balanceCooldown' => 3,
],
],
],
Laravel Horizon 是 Redis 队列的生产级仪表盘,提供实时监控、自动伸缩和指标收集。它通过
1 | balance: auto |
模式自动根据队列压力调整 Worker 数量,在生产环境中非常实用。
五、监控与运维:确保队列健康
队列系统上线后,监控是运维的生命线。核心指标包括队列深度、消费延迟、失败率和 Worker 进程状态。
5.1 关键监控指标
| 指标 | 告警阈值 | 含义 |
|---|---|---|
| 队列深度 | > 10000 或持续增长 | 消费能力不足,需扩容 Worker |
| 消费延迟 | > 60秒(P99) | 任务从入队到处理完成的耗时过长 |
| 失败率 | > 5% | 下游服务异常或代码 Bug |
| Worker 进程数 | < 配置数 | 进程崩溃未恢复 |
| 死信队列深度 | > 0 持续增长 | 需要人工介入排查 |
5.2 自定义监控脚本
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32 <?php
// monitor.php
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$queues = ['emails', 'reports', 'default'];
$alerts = [];
foreach ($queues as $queue) {
$depth = $redis->lLen("queue:{$queue}");
$failed = $redis->lLen("queue:{$queue}:failed");
$reserved = $redis->zCard("queue:{$queue}:reserved");
if ($depth > 5000) {
$alerts[] = "[WARN] Queue {$queue} depth: {$depth}";
}
if ($failed > 0) {
$alerts[] = "[CRITICAL] Queue {$queue} has {$failed} failed jobs";
}
echo sprintf(
"%-15s depth=%-6d reserved=%-4d failed=%-3d\n",
$queue, $depth, $reserved, $failed
);
}
if (!empty($alerts)) {
foreach ($alerts as $alert) {
echo $alert . "\n";
}
exit(1);
}
5.3 RabbitMQ Management API 监控
1
2
3
4
5
6
7
8
9
10
11
12
13 # 获取队列状态
curl -s -u guest:guest \
http://127.0.0.1:15672/api/queues/%2F/tasks.emails | \
jq '{name, messages, messages_ready, messages_unacknowledged, consumers}'
# 输出示例:
# {
# "name": "tasks.emails",
# "messages": 42,
# "messages_ready": 38,
# "messages_unacknowledged": 4,
# "consumers": 3
# }
六、架构演进路线图
从零搭建到生产级队列系统,推荐的演进路线如下:
- 第一阶段(0-1万用户):Redis List + Supervisor Worker,快速上线
- 第二阶段(1-10万用户):Laravel Queue + Horizon,标准化 Job 定义和监控
- 第三阶段(10万+用户):RabbitMQ + 独立微服务 Worker,解耦业务与队列
- 第四阶段(百万级):Kafka + Flink 流处理,支持事件溯源和实时分析
每个阶段的技术选型都应匹配业务规模和团队人力。过早引入 RabbitMQ 会增加运维复杂度,过晚则可能因为 Redis 队列的局限导致数据丢失。关键是在业务增长中找到平衡点。
总结
消息队列是 PHP 应用从同步走向异步的核心基础设施。Redis 队列足够应对轻量场景,RabbitMQ 提供了生产级的可靠投递和路由能力,Laravel Queue 则在框架层提供了优雅的抽象。无论选择哪种方案,都需要关注三个关键维度:可靠性(消息不丢、至少消费一次)、可观测性(队列深度、消费延迟、失败率)和可伸缩性(Worker 扩容、队列分流)。
最终,好的队列架构不是选择了哪个中间件,而是建立了一套从投递到消费、从重试到死信、从监控到告警的完整闭环。掌握了这套闭环,无论业务如何演进,队列系统都能可靠地承载异步任务的压力。
汤不热吧