为什么个人站长会考虑自建消息队列
大多数个人站长的技术栈里其实并不需要消息队列。但总有几个场景会把人逼到这一步:网站要发注册验证邮件,SMTP 一慢,PHP 请求就卡在那里等十几秒,用户以为网站挂了;爬虫或采集任务要抓几千条数据,一次性跑完把服务器 CPU 打满,前台页面直接打不开;还有图片处理,用户上传一张原图后要生成三四种尺寸的缩略图,同步做的话上传接口慢得离谱。这三个场景的共同点是:有些活不需要用户在页面里等着,但它又不适合塞进 crontab——crontab 只有分钟级精度,任务之间无法传递参数,失败了也不知道。
消息队列解决的正是这个问题:把「产生任务」和「执行任务」拆开,Web 进程只负责往队列里扔一条消息就立刻返回,后台的消费进程慢慢处理。生产者不关心消费者是谁、在不在线、处理多久。这个解耦带来的直接好处是前台响应时间稳定,间接好处是任务可以在服务器空闲时批量处理,还能失败重试。
那为什么不装 Kafka、不用 Redis 的 List 凑合?Kafka 对个人站来说太重了,单机跑起来光 Java 进程就要吃 1GB 以上内存,运维复杂度也高。Redis 做队列确实能跑,但它本质上是个内存数据库,消息可靠性依赖 AOF/RDB 持久化策略,一旦配置不当(比如 appendfsync everysec 下断电),未落盘的消息就丢了。RabbitMQ 是个折中:比 Kafka 轻得多,单机 200MB 内存能跑起来,比 Redis List 可靠得多,有真正的消息确认机制(ack)、持久化队列和死信队列,而且它用的是 Erlang,天生为高并发连接设计。
本文假设你是一台 2GB 内存的单机服务器,跑着 Nginx 加 PHP 的网站,想顺手把异步任务这套东西搭起来。下面从安装开始,到实际把 PHP 的一次邮件发送改成异步,完整走一遍,重点讲那些文档里不会写但线上一定会踩的坑。
安装与最小的安全加固
不要用 apt 装 Ubuntu/Debian 源里的 RabbitMQ。源里的版本常年落后,而且依赖的 Erlang 版本也旧,插件的兼容性会出问题。官方推荐用 Cloudsmith 仓库或者直接下 Erlang 与 RabbitMQ 的官方包。不过在个人站的场景下,最省事的其实是 Docker,前提是你已经装好了 Docker 和 Compose。
# docker-compose.yml
services:
rabbitmq:
image: rabbitmq:3.13-management
container_name: rabbitmq
restart: unless-stopped
ports:
- "127.0.0.1:5672:5672" # AMQP 协议,只监听本地
- "127.0.0.1:15672:15672" # 管理界面
environment:
RABBITMQ_DEFAULT_USER: ${RABBITMQ_USER}
RABBITMQ_DEFAULT_PASS: ${RABBITMQ_PASS}
RABBITMQ_DEFAULT_VHOST: /
volumes:
- ./data:/var/lib/rabbitmq
- ./rabbitmq.conf:/etc/rabbitmq/conf.d/10-custom.conf:ro
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 30s
timeout: 10s
retries: 5
ulimits:
nofile:
soft: 65536
hard: 65536三个要点。第一,端口绑定写成 127.0.0.1:15672 而不是 15672。管理界面默认账号密码如果暴露在公网上,会被扫描器在几小时内找到并尝试登录,接着就是被用来当跳板。真要从外面访问,用 SSH 隧道:ssh -L 15672:127.0.0.1:15672 user@server,然后在本地浏览器开 http://localhost:15672。第二,ulimits.nofile 一定要加上。RabbitMQ 每一条连接、每一个队列都要占文件描述符,默认的 1024 在几百个队列时就会开始报错,而且报错信息是莫名其妙的 IO 异常。第三,RABBITMQ_DEFAULT_PASS 不要写在 compose 文件里提交到 Git,用 .env 文件并加进 .gitignore。
如果要写配置文件,rabbitmq.conf 至少要加这几行,用来防住内存爆掉:
# 当内存使用超过 40% 时,阻塞所有生产者连接
vm_memory_high_watermark.relative = 0.4
# 磁盘剩余低于 1GB 时同样阻塞生产者
disk_free_limit.absolute = 1GB
# 消费端预取,防止单个消费者被压垮
# 这一个参数在客户端设置更灵活,配置里只做兜底
channel_max = 2047
heartbeat = 60vm_memory_high_watermark 这个参数是整个 RabbitMQ 运维里最重要的一个。它的行为是:当节点内存使用超过阈值时,RabbitMQ 会阻塞所有发布连接(不是拒绝,是让生产者阻塞住),但消费者不受影响,继续把消息消费掉释放内存。这个设计的意图是「先让消费者把积压消化掉,避免节点被 OOM Killer 杀掉导致消息全丢」。所以如果你的网站发布端突然卡住不响应,第一件事不是怀疑队列坏了,而是去看管理界面的内存水位——很可能队列积压太多,触发保护了。默认值是 0.4,在 2GB 内存的机器上就是约 820MB。这个数字对纯队列场景够用,但如果你在同一台机器上还跑着 MySQL 和 PHP-FPM,就要算清楚这台机器实际能给 RabbitMQ 多少内存,必要时下调到 0.3。
队列、交换机与路由:三个名字别搞混
RabbitMQ 的模型里,生产者从来不是把消息「发到队列」,而是发到交换机(exchange),由交换机根据路由规则决定投递到哪些队列。这个设计第一次接触会觉得很绕,但它带来的灵活性是有价值的。
四种交换机类型里,个人站最常用两种。direct 类型按 routing key 精确匹配,适合「一种任务一个队列」的场景,比如 task.email、task.thumbnail。topic 类型支持通配符,* 匹配一个词、# 匹配零个或多个词,适合「一个消费者订阅一类任务」的场景,比如一个日志消费者用 log.# 就能收全部日志。另外两种 fanout 和 headers 在个人站用得少,前者是广播给所有绑定队列,后者按消息头匹配。
一个典型的坑是 队列没有绑定到交换机,消息发出去就凭空消失了。RabbitMQ 默认不会报错——交换机存在、消息格式合法,它就认为投递成功了,然后因为没有任何队列绑定,消息被直接丢弃。要发现这种情况,你需要在管理界面的 exchange 详情里看 Bindings 那一栏,或者在客户端发布消息时带上 mandatory 标志并注册 return 回调,这样消息无法路由时会被退回给你。
另一个容易被忽略的设置是 队列持久化(durable)与消息持久化(delivery_mode)是两回事,必须同时开才能保证重启不丢消息:
# Python pika 示例:声明一个持久化队列
channel.queue_declare(queue='task.email', durable=True)
# 发布时把消息标记为持久化
channel.basic_publish(
exchange='',
routing_key='task.email',
body=json.dumps(payload),
properties=pika.BasicProperties(
delivery_mode=2, # 2 = persistent,1 = transient
content_type='application/json',
),
)只设 durable=True 而 delivery_mode=1,队列本身在重启后还在,但里面的消息全没了。反过来只设 delivery_mode=2,队列是临时的(默认 durable=False),重启后队列连同消息一起消失。两个都打开之后,代价是每条消息发布时要等一次 fsync,吞吐会下降,但对个人站来说一天几千条消息完全无所谓。如果队列在声明时写了 durable=True,之后用 durable=False 去重复声明同一个队列,RabbitMQ 会直接报 PRECONDITION_FAILED 并关掉这个 channel——因为队列属性不允许中途更改。这个错误在开发环境改代码后特别常见,处理办法是去管理界面把队列删掉重建,或者改个队列名。
消息积压:怎么发现、怎么判断、怎么处理
队列积压是 RabbitMQ 最常见的线上故障。它的表现是:网站前台看起来正常,但用户投诉「密码重置邮件收不到」「上传的图半天不显示缩略图」。这时候去管理界面的 Queues 页面,会看到一个队列的 Ready 数量在持续增长。
三个数字必须搞清楚:Ready 是等待被投递给消费者的消息数;Unacked 是已经投递给消费者但还没收到 ack 的消息数;Total 是两者之和。如果 Ready 高而 Unacked 接近 0,说明消费者根本没在工作——进程挂了、或者连接断了。如果 Unacked 很高,说明消费者拿到了消息但处理不完或者没正确 ack,这时候要继续看消费者的健康度。区分这两个状态能省掉一半的排查时间。
命令行查积压,比开网页更快:
# 查看队列深度(含未确认消息)
rabbitmqctl list_queues name messages messages_ready messages_unacknowledged
# 查看消费者数量,如果为 0 就是没人消费
rabbitmqctl list_queues name consumers
# 查看连接列表,确认消费者进程是否真的连上了
rabbitmqctl list_connections peer_host state
# 用 HTTP API 查(管理插件开启后)
curl -s -u user:pass http://127.0.0.1:15672/api/queues/%2F/task.email | \
python3 -m json.tool | grep -E 'messages|consumers'处理积压有三条路,选择取决于积压的原因。第一,如果是消费者挂了,重启消费者进程即可,但要注意重启后不要让它瞬间拉走全部积压消息然后内存爆掉——这就是 prefetch_count 的作用。
# qos=1 表示同时只给这个消费者 1 条未确认消息
# 处理完 ack 之后才会给下一条
channel.basic_qos(prefetch_count=1)
channel.basic_consume(queue='task.email', on_message_callback=callback)不设 basic_qos 的话,RabbitMQ 会把队列里所有消息尽可能快地推给消费者,消息全部堆在消费者的进程内存里。一个消费者如果手里握着 5 万条邮件任务,进程被杀掉后这 5 万条消息会因为没有 ack 而退回队列,然后被重新推给新的消费者,如此循环,队列永远清不掉。这就是为什么 prefetch_count 必须设,而且对个人站来说设成 1 到 5 就够了。
第二,如果是消费速度跟不上产生速度,那要么加消费者进程(RabbitMQ 天然支持多个消费者竞争同一个队列,消息会被轮询分配),要么优化消费逻辑。第三,如果积压的消息已经过期没有意义了(比如三天前的验证邮件),直接清空队列比慢慢消费更快:管理界面队列详情页有 Purge 按钮,命令行是 rabbitmqctl purge_queue task.email。但清空操作的返回是「被删除的消息条数」,这个动作不可逆,执行前要确认这个队列的消息确实可以丢。
死信队列:让失败的任务有地方可查
消息在三种情况下会变成「死信」:被消费者 basic_nack 且 requeue=false 拒绝;消息在队列里存活时间超过设置的 TTL;队列长度超过限制被挤出的旧消息。这些消息如果没有任何处理,就直接消失了,你永远不知道哪封邮件没发出去。
正确做法是给业务队列配一个死信交换机和死信队列,所有进入死信状态的消息自动转到那里,之后人工检查或者写个脚本统一重试:
# 1. 声明死信交换机与队列
channel.exchange_declare(exchange='dlx', exchange_type='direct', durable=True)
channel.queue_declare(queue='dlx.email', durable=True)
channel.queue_bind(queue='dlx.email', exchange='dlx', routing_key='email.failed')
# 2. 业务队列声明时指定死信参数
channel.queue_declare(
queue='task.email',
durable=True,
arguments={
'x-dead-letter-exchange': 'dlx',
'x-dead-letter-routing-key': 'email.failed',
'x-message-ttl': 86400000, # 24 小时没被消费就转死信
'x-max-length': 10000, # 队列最多 1 万条,超出从头丢
},
)注意 x-message-ttl 的单位是毫秒,86400000 才是 24 小时,写 86400 就变成了 86 秒,这个错误非常隐蔽——队列里的消息会在你还没反应过来时全部变成死信。另外队列级别的 TTL 是「从入队开始计时」,而给单条消息设置 TTL 是「排在队头才开始计时」,后者在积压时行为会很怪异(因为队头消息的 TTL 到期前,后面所有消息的 TTL 都不开始计算),所以个人站直接用队列级 TTL 就够了。
死信队列里的消息会保留原始消息的所有属性,包括 x-death 头,里面记录了它来自哪个队列、死亡原因是什么:
{
"x-death": [{
"count": 3,
"reason": "expired",
"queue": "task.email",
"exchange": "",
"routing-keys": ["task.email"],
"time": 1759100000
}]
}reason 字段是关键,expired 表示 TTL 到期,rejected 表示被消费者拒绝,maxlen 表示被队列长度限制挤出。看到 rejected 且 count 很大,说明消费逻辑有 bug 在反复失败;看到 expired 一片,说明消费者长时间不在线。
把它接进现有的 PHP 网站
最后落到实际场景:把密码重置邮件的发送从同步改成异步。原来 PHP 里是 mail() 或者 PHPMailer 直接调 SMTP,用户点「发送重置链接」后要等 SMTP 服务器响应,慢的时候 5 到 15 秒。改造后,PHP 只负责把任务扔进队列,立刻返回「邮件已发送,请查收」。
// composer require php-amqplib/php-amqplib
require_once __DIR__ . '/vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
function enqueueEmail(string $to, string $subject, string $html): void {
$conn = new AMQPStreamConnection(
'127.0.0.1', 5672, getenv('RABBITMQ_USER'), getenv('RABBITMQ_PASS')
);
$ch = $conn->channel();
$ch->queue_declare('task.email', false, true, false, false);
$body = json_encode([
'to' => $to,
'subject' => $subject,
'html' => $html,
'ts' => time(),
], JSON_UNESCAPED_UNICODE);
$msg = new AMQPMessage($body, [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
'content_type' => 'application/json',
]);
$ch->basic_publish($msg, '', 'task.email'); // 空 exchange = 默认交换机
$ch->close();
$conn->close();
}有个性能细节:每次请求都新建连接和 channel 是很浪费的。AMQP 建连接要握手,一次要几毫秒到几十毫秒。如果 PHP 跑在 php-fpm 下,可以改用 AMQPLazyConnection,它会延迟到真正发送消息时才建连接,对于「大部分请求不发消息」的页面能省掉开销。真正的连接池在 PHP-FPM 的模型下不好做,因为进程是复用的但请求之间会清理资源,务实的选择是接受每次建连接的成本,或者干脆在 Web 端不直接连 RabbitMQ,而是先写进本地的一张小表,用一个小脚本批量搬运——但这又牺牲了实时性,看你的量级。
消费者的写法要点:
$ch->basic_qos(null, 1, null);
$ch->basic_consume('task.email', '', false, false, false, false,
function ($msg) {
try {
sendMail(json_decode($msg->body, true));
$msg->ack(); // 成功才确认
} catch (Throwable $e) {
error_log('send mail failed: ' . $e->getMessage());
$msg->nack(false, false); // requeue=false → 进死信队列
}
}
);
while ($ch->is_consuming()) {
$ch->wait();
}必须保证 ack 只在真正成功之后调用。如果代码是先 ack 再发邮件,那么在发邮件过程中进程崩溃,这条消息就永久丢失了。反过来,如果异常时用了 nack(false, true)(requeue=true),一条注定失败的消息(比如收件人地址格式非法)会被无限次重投,把 CPU 吃满。正确的策略是:可预期的业务失败直接 nack 进死信队列,暂时性的故障(SMTP 连不上)才 requeue,而且最好加个重试计数,超过三次就送死信。
让消费者在服务器重启后自己起来
用手动 php consumer.php 启动的消费者,一旦服务器重启、SSH 断开或者进程崩溃,就没了。要用 systemd 把它管起来:
# /etc/systemd/system/mq-email-consumer.service
[Unit]
Description=RabbitMQ email consumer
After=network-online.target docker.service
Wants=network-online.target
[Service]
Type=simple
User=www-data
WorkingDirectory=/var/www/site
EnvironmentFile=/etc/mq-consumer.env
ExecStart=/usr/bin/php /var/www/site/consumer_email.php
Restart=always
RestartSec=10
StandardOutput=append:/var/log/mq-email-consumer.log
StandardError=append:/var/log/mq-email-consumer.log
[Install]
WantedBy=multi-user.targetRestart=always 加上 RestartSec=10 是必须的。如果 RabbitMQ 容器启动比消费者慢,消费者连不上会退出,systemd 每 10 秒重试一次,等 RabbitMQ 就绪后自动恢复。不要用 Restart=on-failure,因为消费者进程正常退出(比如队列被删导致 channel 关闭)同样需要重启。另外 StandardOutput=append: 这个写法要求 systemd 236 以上,如果你的系统旧,改成 StandardOutput=journal 然后 journalctl -u mq-email-consumer -f 看日志。
加日志轮转,否则这个日志文件会一直涨:
# /etc/logrotate.d/mq-consumer
/var/log/mq-email-consumer.log {
daily
rotate 7
compress
delaycompress
missingok
notifempty
copytruncate
}copytruncate 在这里是必要的,因为 systemd 保持了对这个文件描述符的持有,普通的 rename 方式轮转后 systemd 还会往已经改名的旧文件里写。用 copytruncate 会先复制再清空原文件,代价是复制瞬间可能丢极少量日志,但对消费者日志完全可接受。
内存和连接的日常监控
RabbitMQ 本身有 Prometheus 插件,但个人站没必要搭一整套监控,用一条 cron 脚本检查关键指标、超阈值时给自己发个邮件或从别的通道告警就够了:
#!/bin/bash
# /usr/local/bin/mq-check.sh —— 每 10 分钟跑一次
API="http://127.0.0.1:15672/api"
AUTH="user:pass"
ALERT_THRESHOLD=5000 # 积压告警阈值
# 1. 队列积压
for q in task.email task.thumbnail; do
n=$(curl -s -u "$AUTH" "$API/queues/%2F/$q" | \
python3 -c "import sys,json;print(json.load(sys.stdin)['messages'])" 2>/dev/null)
if [ -n "$n" ] && [ "$n" -gt "$ALERT_THRESHOLD" ]; then
echo "ALERT: queue $q backlog=$n"
fi
done
# 2. 消费者数量为 0 告警
cons=$(curl -s -u "$AUTH" "$API/queues/%2F/task.email" | \
python3 -c "import sys,json;print(json.load(sys.stdin)['consumers'])" 2>/dev/null)
[ "$cons" = "0" ] && echo "ALERT: task.email has no consumer"
# 3. 内存水位(VmRSS 超过 800MB 告警)
mem=$(curl -s -u "$AUTH" "$API/nodes" | \
python3 -c "import sys,json;print(json.load(sys.stdin)[0]['mem_used'])" 2>/dev/null)
[ -n "$mem" ] && [ "$mem" -gt 838860800 ] && echo "ALERT: rabbit mem=${mem}B"
exit 0消费者数为 0 这条告警比积压告警更重要,因为它能在用户还没察觉之前就发现问题。积压到 5000 条可能已经过去半小时了,而消费者掉了的瞬间就能知道。
最后提醒两个容易忽略的运维点。一是 RabbitMQ 的数据目录一定要放在有空间的磁盘上,默认 /var/lib/rabbitmq 如果在系统盘而系统盘只有 20GB,永久队列的元数据和持久化消息会把盘写满,写满之后 RabbitMQ 会拒绝所有消息(这就是 disk_free_limit 的作用)。二是 定期做一次恢复演练:停掉容器、删掉队列里的消息、重启,确认生产者能重连、消费者能重新注册、死信队列里的消息还在。docker compose restart rabbitmq 之后等 30 秒,用 rabbitmqctl list_queues name durable messages 检查队列是否都回来了、durable 那一列全是 true。没做过恢复演练的备份等于没有备份,这句话对消息队列同样成立。