Skip to content

RabbitMQ 最佳实践

本节总结 RabbitMQ 在 PHP 项目中的最佳实践,涵盖消息持久化与可靠投递、流量控制与背压管理、集群部署策略、监控告警体系以及消息幂等性处理等核心内容,帮助团队构建可靠、高效的消息队列系统。

前置知识

阅读本节前,建议先了解:RabbitMQ 基础PHP 客户端操作

消息持久化与可靠投递

可靠性三要素

确保消息不丢失需要以下三要素同时满足:

1. Exchange 持久化(durable: true)
2. Queue 持久化(durable: true)
3. Message 持久化(delivery_mode: 2)

缺一不可!
php
<?php

declare(strict_types=1);

namespace App\MessageQueue;

use PhpAmqpLib\Channel\AMQPChannel;
use PhpAmqpLib\Message\AMQPMessage;

class ReliablePublisher
{
    public function __construct(
        private readonly AMQPChannel $channel
    ) {}

    /**
     * 发送可靠消息(持久化 + Publisher Confirm)
     */
    public function publishReliable(
        string $exchange,
        string $routingKey,
        array $data
    ): bool {
        // 1. 开启 Publisher Confirm
        $this->channel->confirm_select();

        // 2. 设置 mandatory=true(消息无法路由时返回给 Producer)
        $mandatory = true;

        // 3. 构造持久化消息
        $message = new AMQPMessage(
            body: json_encode($data),
            properties: [
                'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
                'content_type'  => 'application/json',
                'message_id'    => uniqid('msg_', true),
                'priority'      => 0,
                'timestamp'     => time(),
            ]
        );

        // 4. 发送消息
        try {
            $this->channel->basic_publish(
                msg: $message,
                exchange: $exchange,
                routing_key: $routingKey,
                mandatory: $mandatory,
            );

            // 5. 等待 Broker 确认
            $this->channel->wait_for_pending_acks(timeout: 5.0);

            return true;
        } catch (\Throwable $e) {
            // 处理失败:重试或记录
            return false;
        }
    }
}

消息可靠性保障策略

Producer 端保障:
├── Publisher Confirm(确认消息到达 Exchange)
├── mandatory 标志(消息无法路由时返回)
└── Return Listener(处理无法路由的消息)

Broker 端保障:
├── Exchange 持久化
├── Queue 持久化
├── Message 持久化
└── 镜像队列 / Quorum Queue(集群级别)

Consumer 端保障:
├── 手动 ACK(处理成功后确认)
├── NACK + 重试(处理失败后重试)
└── 死信队列(最终兜底)

流量控制

Prefetch Count

prefetch_count 是 RabbitMQ 流量控制的核心参数,控制每次向 Consumer 推送多少条未确认消息。

php
<?php

declare(strict_types=1);

namespace App\MessageQueue;

use PhpAmqpLib\Channel\AMQPChannel;
use PhpAmqpLib\Message\AMQPMessage;

class FlowControlConsumer
{
    public function __construct(
        private readonly AMQPChannel $channel
    ) {}

    /**
     * 根据消费者处理能力设置 prefetch_count
     *
     * 原则:
     * - IO 密集型(如调用外部 API):prefetch_count = 10~20
     * - CPU 密集型(如数据处理):prefetch_count = 1~5
     * - 快速处理(如日志写入):prefetch_count = 100~500
     */
    public function consumeWithPrefetch(
        string $queue,
        callable $handler,
        int $prefetchCount = 10
    ): void {
        // 设置 QoS(Quality of Service)
        $this->channel->basic_qos(
            prefetch_size: 0,           // 不限制消息大小
            prefetch_count: $prefetchCount,
            global: false,             // 仅对当前 Consumer 生效
        );

        $this->channel->basic_consume(
            queue: $queue,
            no_ack: false,
            callback: function (AMQPMessage $message) use ($handler): void {
                try {
                    $data = json_decode($message->body, true);
                    $handler($data);
                    $message->ack();
                } catch (\Throwable $e) {
                    $message->nack(requeue: true);
                }
            },
        );

        while ($this->channel->is_consuming()) {
            $this->channel->wait(timeout: 30);
        }
    }
}

消费者限流

php
<?php

declare(strict_types=1);

namespace App\MessageQueue;

/**
 * 令牌桶限流器(控制消费速率)
 */
class RateLimiter
{
    private int $maxRate;
    private int $tokens;
    private int $lastRefillTime;

    public function __construct(int $maxRate = 1000)
    {
        $this->maxRate = $maxRate;
        $this->tokens = $maxRate;
        $this->lastRefillTime = microtime(true);
    }

    /**
     * 尝试获取一个令牌
     */
    public function allow(): bool
    {
        $this->refill();

        if ($this->tokens > 0) {
            $this->tokens--;
            return true;
        }

        return false;
    }

    /**
     * 补充令牌
     */
    private function refill(): void
    {
        $now = microtime(true);
        $elapsed = $now - $this->lastRefillTime;

        $tokensToAdd = (int) ($elapsed * $this->maxRate);
        if ($tokensToAdd > 0) {
            $this->tokens = min($this->maxRate, $this->tokens + $tokensToAdd);
            $this->lastRefillTime = $now;
        }
    }
}

集群部署

集群架构

RabbitMQ 集群架构

Client → Node1 (Master)
         Node2 (Slave)
         Node3 (Slave)

节点角色:
- 所有节点都可以接受连接
- 队列数据默认在声明它的节点上
- 使用镜像队列(Classic Mirrored Queue)或 Quorum Queue 实现数据复制

Quorum Queue(推荐)

Quorum Queue 是 RabbitMQ 3.8+ 引入的新型队列,基于 Raft 共识算法实现高可用。

php
<?php

declare(strict_types=1);

// Quorum Queue 声明(替代传统镜像队列)
$this->channel->queue_declare(
    queue: 'order.queue',
    durable: true,
    arguments: [
        'x-queue-type' => 'quorum',           // Quorum Queue
    ],
);

// 对比传统队列
// Quorum Queue vs Classic Queue
// - 数据一致性更强(Raft 共识)
// - 故障恢复更快
// - 消息不丢失保证
// - 性能略低于 Classic Queue

队列类型选择

  • Classic Queue(+ 镜像):适合低延迟、高吞吐量场景
  • Quorum Queue:适合数据可靠性要求高的场景(推荐)
  • Stream Queue:适合日志、审计等大流量追加场景

集群管理命令

bash
# 查看集群状态
rabbitmqctl cluster_status

# 查看节点列表
rabbitmqctl cluster_status | grep "running_nodes"

# 添加节点到集群
rabbitmqctl stop_app
rabbitmqctl reset
rabbitmqctl join_cluster rabbit@node1
rabbitmqctl start_app

# 从集群移除节点
rabbitmqctl stop_app
rabbitmqctl reset
rabbitmqctl start_app

# 设置镜像队列策略
rabbitmqctl set_policy ha-all "^" '{"ha-mode":"all","ha-sync-mode":"automatic"}' --apply-to queues

# 查看 Web 管理界面
# http://node1:15672

监控与告警

关键监控指标

php
<?php

declare(strict_types=1);

namespace App\Monitor;

class RabbitMQMonitor
{
    /**
     * 监控指标清单
     *
     * 队列级别:
     * - messages_ready:等待被消费的消息数
     * - messages_unacked:已投递但未确认的消息数
     * - consumers:当前活跃消费者数
     * - message_stats.publish_details.rate:消息发布速率
     * - message_stats.deliver_get_rate:消息消费速率
     *
     * 节点级别:
     * - mem_used:内存使用量
     * - disk_free:磁盘剩余空间
     * - socket_used:TCP 连接数
     * - process_used:进程使用率
     *
     * 集群级别:
     * - partitions:网络分区数(> 0 为异常)
     * - running_nodes:运行中节点数
     */
    public function getQueueMetrics(): array
    {
        // 通过 RabbitMQ Management API 获取
        $apiUrl = 'http://127.0.0.1:15672/api/queues';

        $ch = curl_init();
        curl_setopt_array($ch, [
            CURLOPT_URL     => $apiUrl,
            CURLOPT_USERPWD => 'admin:admin123',
            CURLOPT_RETURNTRANSFER => true,
        ]);

        $response = curl_exec($ch);
        $queues = json_decode($response, true);

        $metrics = [];
        foreach ($queues as $queue) {
            $metrics[] = [
                'name'         => $queue['name'],
                'messages'     => $queue['messages'] ?? 0,
                'messages_ready'  => $queue['messages_ready'] ?? 0,
                'messages_unacked'=> $queue['messages_unacked'] ?? 0,
                'consumers'    => $queue['consumers'] ?? 0,
            ];
        }

        return $metrics;
    }

    /**
     * 健康检查
     */
    public function healthCheck(): array
    {
        $apiUrl = 'http://127.0.0.1:15672/api/overview';
        $ch = curl_init();
        curl_setopt_array($ch, [
            CURLOPT_URL     => $apiUrl,
            CURLOPT_USERPWD => 'admin:admin123',
            CURLOPT_RETURNTRANSFER => true,
        ]);

        $response = curl_exec($ch);
        $overview = json_decode($response, true);

        $nodeCount = $overview['rabbitmq_version'] ? 1 : 0;
        $queueTotals = $overview['queue_totals'] ?? [];
        $messageStats = $overview['message_stats'] ?? [];

        return [
            'cluster_name'       => $overview['cluster_name'] ?? 'unknown',
            'messages_ready'     => $queueTotals['messages_ready'] ?? 0,
            'messages_unacked'   => $queueTotals['messages_unacked'] ?? 0,
            'object_totals'      => $overview['object_totals'] ?? [],
            'status'             => 'healthy',
        ];
    }
}

告警规则

告警项阈值级别处理方式
队列消息堆积> 10000Warning检查消费者是否正常
队列消息堆积> 100000Critical扩容消费者或降级
消费者为 0队列有消息Warning检查消费者进程
节点内存 > 80%-Warning检查是否有内存泄漏
磁盘空间 < 20%-Critical清理日志或扩容磁盘
网络分区> 0Critical检查网络,可能需要人工干预

消息幂等性

什么是幂等性

幂等性(Idempotency)是指一个操作执行一次和执行多次的效果相同。在消息队列场景中,消费者可能会收到重复消息(网络重发、消费者重启等),因此必须实现幂等处理。

php
<?php

declare(strict_types=1);

namespace App\MessageQueue;

use Redis;

/**
 * 幂等性消费者
 * 使用 message_id 去重,防止重复处理
 */
class IdempotentConsumer
{
    public function __construct(
        private readonly Redis $redis
    ) {}

    /**
     * 幂等性处理消息
     * 使用 Redis SETNX 实现消息去重
     */
    public function handleWithIdempotent(string $messageId, string $queueName, callable $handler): bool
    {
        $dedupKey = "msg_dedup:{$queueName}:{$messageId}";

        // 尝试设置去重标记(仅当 key 不存在时成功)
        $isFirstTime = $this->redis->set($dedupKey, '1', ['NX', 'EX' => 86400]);

        if (!$isFirstTime) {
            // 消息已被处理过,跳过
            return true;
        }

        try {
            $result = $handler();
            if (!$result) {
                // 处理失败,移除去重标记(允许重试)
                $this->redis->del($dedupKey);
            }
            return $result;
        } catch (\Throwable $e) {
            $this->redis->del($dedupKey);
            throw $e;
        }
    }

    /**
     * 基于业务唯一标识的幂等处理
     * 例如订单号、支付流水号等
     */
    public function handleByBusinessId(
        string $businessType,
        string $businessId,
        callable $handler
    ): mixed {
        $dedupKey = "biz_dedup:{$businessType}:{$businessId}";

        $isFirstTime = $this->redis->set($dedupKey, 'processing', ['NX', 'EX' => 3600]);

        if (!$isFirstTime) {
            // 已处理,直接返回
            return null;
        }

        try {
            $result = $handler();
            // 标记为已完成
            $this->redis->setex($dedupKey, 86400, 'completed');
            return $result;
        } catch (\Throwable $e) {
            $this->redis->del($dedupKey);
            throw $e;
        }
    }
}

幂等性策略对比

策略实现方式适用场景
message_id 去重Redis SETNX + TTL消息有唯一 ID
业务 ID 去重Redis SETNX + 业务唯一标识订单号、支付号等
数据库唯一约束UNIQUE INDEX + INSERT IGNORE最终写入数据库的场景
状态机订单状态从待处理→已处理有明确状态流转的场景
版本号(CAS)乐观锁(version 字段)更新操作

注意事项

  1. Exchange/Queue 命名规范:使用点分隔的层级命名,如 order.created.exchange
  2. 消息版本控制:在消息体中包含版本号,便于升级兼容
  3. 消息大小:单条消息建议不超过 1MB,超大消息使用外部存储 + 消息中传递引用
  4. 消费者组:同一消费者组的多个实例共享负载,使用相同的 queue 名
  5. 连接恢复:实现自动重连机制,应对网络抖动

最佳实践总结

  1. 可靠性:Exchange/Queue/Message 三重持久化 + Publisher Confirm + 手动 ACK
  2. 流量控制:合理设置 prefetch_count,避免消费者过载
  3. 高可用:使用 Quorum Queue 或镜像队列,部署多节点集群
  4. 幂等性:所有消费者实现幂等处理,使用 Redis 去重或数据库唯一约束
  5. 监控:监控队列堆积、消费者数量、消息速率,设置合理告警阈值
  6. 消息追踪:使用 message_id 和 tracing 实现全链路追踪

下一节

继续学习:Kafka 基础

参考链接