本文目录导读:

在PHP中保证消息不丢失(即消息的可靠性传递),需要从生产端、消费端、中间件三个层面综合处理,以下是我总结的完整解决方案:
生产端可靠性(确保消息发出去)
确认机制(ACK)
// RabbitMQ 生产者确认
$channel->confirm_select(); // 开启确认模式
$channel->set_ack_handler(function() {
echo "消息确认送达";
});
$channel->set_nack_handler(function() {
echo "消息发送失败,需要重试";
});
$channel->basic_publish($msg, $exchange, $routingKey);
$channel->wait_for_pending_acks(); // 等待确认
持久化存储
// RabbitMQ 持久化配置
$channel->exchange_declare($exchange, 'direct',
false, // passive
true, // durable 交换机持久化
false
);
$channel->queue_declare($queue,
false,
true, // durable 队列持久化
false,
false
);
// 消息持久化
$msg = new AMQPMessage($data, [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT
]);
本地缓存 + 定时重发
class ReliableProducer {
private $pendingMessages = [];
public function send($message) {
// 1. 先写入本地缓存
$this->pendingMessages[] = $message;
// 2. 尝试发送
try {
$this->publish($message);
// 3. 发送成功移除
array_shift($this->pendingMessages);
} catch (Exception $e) {
// 记录日志,等待重试
}
}
public function retryPending() {
foreach ($this->pendingMessages as $msg) {
$this->send($msg);
}
}
}
消费端可靠性(确保消息被处理)
手动ACK
// RabbitMQ 手动确认
$channel->basic_consume($queue, '', false,
false, // no_ack = false,关闭自动确认
false, false,
function($msg) {
try {
// 处理业务逻辑
$this->process($msg->body);
// 处理成功后确认
$msg->delivery_info['channel']->basic_ack(
$msg->delivery_info['delivery_tag']
);
} catch (Exception $e) {
// 处理失败,重新入队
$msg->delivery_info['channel']->basic_nack(
$msg->delivery_info['delivery_tag'],
false, // multiple
true // requeue
);
}
}
);
Redis队列的ACK机制
class RedisReliableQueue {
private $redis;
private $processingKey = 'queue:processing';
public function processMessage() {
// 1. 从队列取出消息,放入处理中队列
$message = $this->redis->lpop('queue:pending');
if ($message) {
$this->redis->rpush($this->processingKey, $message);
try {
// 2. 执行业务处理
$this->doBusiness($message);
// 3. 处理成功,从处理中队列删除
$this->redis->lrem($this->processingKey, 0, $message);
} catch (Exception $e) {
// 4. 处理失败,重试或放入死信队列
$this->retryOrDeadLetter($message);
}
}
}
}
中间件配置
RabbitMQ配置
// 集群配置
$connection = new AMQPStreamConnection(
'host1', 5672, 'user', 'pass',
'/', false, 'AMQPLAIN', null, 'en_US',
10, // connection timeout
30, // read timeout
10, // write timeout
60 // heartbeat
);
// 队列参数详解
$arguments = [
'x-ha-policy' => 'all', // 镜像队列
'x-max-length' => 10000, // 队列长度限制
'x-dead-letter-exchange' => 'dlx', // 死信交换机
'x-dead-letter-routing-key' => 'dead' // 死信路由键
];
Kafka配置
// Kafka 生产者配置
$conf = new RdKafka\Conf();
$conf->set('producer.acks', 'all'); // 所有副本确认
$conf->set('producer.retries', '3'); // 重试次数
$conf->set('producer.linger.ms', '100'); // 批量等待时间
$conf->set('producer.batch.size', '16384'); // 批量大小
// Kafka 消费者配置
$conf->set('enable.auto.commit', 'false'); // 关闭自动提交
$conf->set('auto.offset.reset', 'earliest'); // 从头开始消费
业务层面的兜底方案
消息追踪表
CREATE TABLE message_tracking (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
message_id VARCHAR(64) UNIQUE,
content TEXT,
status ENUM('pending', 'sent', 'processed', 'failed') DEFAULT 'pending',
retry_count INT DEFAULT 0,
max_retry INT DEFAULT 3,
created_at TIMESTAMP,
updated_at TIMESTAMP,
INDEX idx_status (status),
INDEX idx_retry (retry_count)
);
幂等性处理
class IdempotentConsumer {
private $redis;
public function handle($messageId, $callback) {
// 1. 检查是否已处理过
$key = "processed:{$messageId}";
if ($this->redis->set($key, 1, ['NX', 'EX' => 3600])) {
// 2. 首次处理
$callback();
$this->redis->set("success:{$messageId}", 1);
}
}
}
定时补偿任务
class MessageReconciler {
public function reconcile() {
// 查询超时未确认的消息
$timedOut = $this->mysql->query(
"SELECT * FROM message_tracking
WHERE status = 'pending'
AND created_at < NOW() - INTERVAL 5 MINUTE"
);
foreach ($timedOut as $msg) {
// 检查是否真的被处理
if (!$this->isProcessed($msg['message_id'])) {
// 重新发送
$this->producer->resend($msg['content']);
// 更新重试次数
$this->mysql->execute(
"UPDATE message_tracking
SET retry_count = retry_count + 1
WHERE message_id = ?"
);
}
}
}
}
最佳实践建议
- 消息ID唯一性:每条消息包含唯一ID,用于追踪
- 超时重试:设置合理的超时时间,超时自动重试
- 死信队列:处理失败的消息进入死信队列,人工处理
- 监控告警:设置消息积压、消费延迟等监控告警
- 日志记录:完整记录消息生命周期日志
PHP保证消息不丢失需要:
- 生产端:使用ACK机制 + 局部存储 + 失败重试
- 消费端:手动ACK + 幂等处理 + 死信队列
- 中间件:消息持久化 + 镜像队列 + 合理配置
- 业务层:消息追踪表 + 定时补偿 + 监控告警
最安全的方案是多重保障,单靠某一层无法完全避免消息丢失。