本文目录导读:

- 核心原则
- 方案一:基于数据库主键/唯一索引(最常用、最可靠)
- 方案二:基于Redis的幂等锁(高并发场景)
- 方案三:业务表状态标记(适合多步骤流程)
- 方案四:专用幂等中间件/表(通用方案)
- 方案五:消息队列自带特性(Kafka、RabbitMQ等)
- 最佳实践建议
- 示例:完整PHP消费代码(结合Redis + 数据库)
在PHP项目中实现消息幂等性(防止重复消费)的核心思路是 保证同一操作无论执行多少次,其结果都一致,针对RabbitMQ、Kafka、Redis队列或本地任务队列等场景,以下是几种可靠且实用的解决方案:
核心原则
- 唯一标识:每条消息必须有一个全局唯一的ID(如UUID、雪花算法ID)。
- 状态记录:一旦消息被成功消费,就记录该ID及其状态。
- 拦截重放:后续收到相同ID的消息时,直接跳过处理。
基于数据库主键/唯一索引(最常用、最可靠)
适用场景:消息处理结果最终要写入数据库(订单、支付、积分等)。
步骤:
-
生成唯一ID:生产者发送消息时生成
message_id(如order_12345)。 -
表结构设计:在业务表中增加
message_id字段,并设置唯一索引。CREATE TABLE user_orders ( id INT AUTO_INCREMENT PRIMARY KEY, order_id VARCHAR(50) NOT NULL, message_id VARCHAR(64) UNIQUE, -- 幂等键 status TINYINT DEFAULT 0, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ); -
消费逻辑:
// 消费者处理消息 function handleMessage(array $message): void { $messageId = $message['message_id']; $orderData = $message['data']; try { // 尝试插入乐此不疲,利用唯一索引防重 DB::insert('INSERT INTO user_orders (order_id, message_id, status) VALUES (?, ?, 1) ON DUPLICATE KEY UPDATE id=id', [$orderData['order_id'], $messageId]); // 如果上面执行成功(影响行数>0),说明是新消息 // 继续执行业务逻辑... } catch (\PDOException $e) { if ($e->getCode() === '23000') { // 违反唯一索引 // 重复消息,直接跳过 return; } throw $e; } }
优点:无需额外中间件,强一致性,数据库自带事务保证。
缺点:对数据库有写压力;不适合非数据库场景(如发短信、调用第三方API)。
基于Redis的幂等锁(高并发场景)
适用场景:消息处理不涉及数据库,或数据库不支持唯一键(如MongoDB),或要求极低延迟。
核心:利用Redis的 SET NX EX 命令实现原子性锁。
步骤:
-
生成唯一ID:同上。
-
消费时加锁:
function handleMessage(array $message): void { $messageId = $message['message_id']; $redisKey = "idempotent:{$messageId}"; // 尝试加锁,有效期5分钟(根据实际处理时间调整) // NX:只有key不存在时才设置,EX:过期时间 $locked = Redis::set($redisKey, '1', 'NX', 'EX', 300); if (!$locked) { // 已经处理过,直接ACK并返回 Log::info("重复消息跳过: {$messageId}"); return; } try { // 执行实际业务逻辑(小心:加锁成功不代表一定执行成功) $this->processBusiness($message); // 业务成功:锁会自动过期,无需额外处理 // 如果需要长期记录,可单独存入集合:Redis::sAdd('processed_messages', $messageId); } catch (\Exception $e) { // 业务失败:必须手动删除锁,否则消息会丢失 Redis::del($redisKey); throw $e; // 触发消息重试机制 } }
关键点:
- 锁必须带过期时间,防止死锁。
- 业务失败时要主动删锁,否则消息永远无法重试。
- 锁时间要大于业务最大耗时,否则可能出现重复处理(阈值问题)。
缺点:Redis宕机或网络分区可能导致幂等失效(概率极低);需要额外维护Redis。
业务表状态标记(适合多步骤流程)
适用场景:消息可能重复触发同一状态机,如“支付成功”后再次收到“支付成功”。
步骤:
-
表中记录当前状态:如
order.status = 'paid'。 -
消费时判断状态:
function handlePaymentSuccess($orderId, $messageId) { $order = DB::table('orders')->where('id', $orderId)->first(); // 如果已经支付,直接跳过 if ($order->status === 'paid') { Log::info("订单已支付,忽略重复消息: {$orderId}"); return; } // 使用乐观锁或行级锁更新状态 DB::beginTransaction(); try { $updated = DB::table('orders') ->where('id', $orderId) ->where('status', 'pending') // 条件更新 ->update(['status' => 'paid', 'message_id' => $messageId]); if ($updated === 0) { // 并发下被其他线程更新,视为重复 DB::rollBack(); return; } // 完成后续操作(积分、短信等) DB::commit(); } catch (\Exception $e) { DB::rollBack(); throw $e; } }
优点:与业务逻辑深度结合,不依赖额外中间件。
缺点:对状态机的设计要求较高,不适合无状态的简单操作。
专用幂等中间件/表(通用方案)
适用场景:需要统一管理多个业务场景的幂等性,或者业务无法修改表结构。
步骤:
-
建一张幂等表:
CREATE TABLE idempotent_records ( message_id VARCHAR(64) PRIMARY KEY, -- 唯一主键 status TINYINT NOT NULL DEFAULT 0, -- 0:处理中,1:成功,2:失败 created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, expired_at TIMESTAMP -- 定期清理 ); -
消费流程:
function handle(array $message): void { $messageId = $message['message_id']; // 1. 先查询处理结果 $record = DB::table('idempotent_records')-> where('message_id', $messageId)->first(); if ($record && $record->status === 1) { Log::info("消息已处理成功: {$messageId}"); return; } // 2. 标记处理中(防止并发) DB::table('idempotent_records')->insertOrUpdate([ 'message_id' => $messageId, 'status' => 0, ]); try { // 3. 执行业务逻辑 doBusiness($message['data']); // 4. 标记成功 DB::table('idempotent_records') ->where('message_id', $messageId) ->update(['status' => 1]); } catch (\Exception $e) { // 标记失败(或删除记录,以便重试) DB::table('idempotent_records') ->where('message_id', $messageId) ->update(['status' => 2]); throw $e; } }
管理建议:
- 定期清理过期记录(如
WHERE expired_at < NOW())。 message_id设计时最好包含业务前缀(如order_xxx、pay_xxx),方便区分和清理。
消息队列自带特性(Kafka、RabbitMQ等)
Kafka
- 启用幂等生产:
enable.idempotence=true,防止生产者重复发送。 - 消费者去重:需配合 雪花算法ID + 外部存储 实现。
RabbitMQ
- 手动ACK + 幂等表:消费者处理成功后手动ACK,再配合方案一或四防重。
- 死信队列:处理失败的消息进入死信队列,由专门消费者重试,重试时同样需幂等。
最佳实践建议
-
优先使用数据库唯一索引(方案一):
- 简单、可靠、事务一致性好。
- 适合95%的业务场景。
-
高并发且数据不落库(如缓存更新、短信验证码)使用Redis锁(方案二):
注意锁时间、异常删除锁。
-
通用去重中间件(方案四):
- 适合微服务架构,多个服务共享幂等逻辑。
- 可配合Redis做缓存加速(先查Redis,再查DB)。
-
生产者侧也需做一定控制:
减少生产重复消息的概率(如使用消息队列的事务机制)。
-
异常处理:
- 如果业务执行失败(如第三方接口超时),不要直接ACK,让消息重试(最多重试3次)。
- 重试时幂等表应保留失败记录或删除锁,允许再次尝试。
-
定期清理:
幂等记录设置过期时间(如72小时),避免表数据无限增长。
示例:完整PHP消费代码(结合Redis + 数据库)
class OrderConsumer {
public function consume($message) {
$messageId = $message['message_id'];
$orderId = $message['order_id'];
// Step1: Redis幂等检查
$redisKey = "order:{$orderId}";
$lock = Redis::set($redisKey, $messageId, 'NX', 'EX', 60);
if (!$lock) {
// 正在处理或已处理
$processedMsgId = Redis::get($redisKey);
if ($processedMsgId === $messageId) {
Log::info("相同消息已处理: {$messageId}");
return true; // ACK
}
// 不同消息但相同业务ID(需按业务判断)
return true; // 跳过
}
try {
// Step2: 数据库幂等(更可靠)
DB::beginTransaction();
$updated = DB::table('orders')
->where('order_id', $orderId)
->where('status', 'pending')
->update(['status' => 'paid']);
if ($updated === 0) {
DB::rollBack();
return true; // 已处理
}
// 其他业务...
DB::commit();
// 设置Redis为永久标记(或较长过期时间)
Redis::set($redisKey, $messageId);
Redis::expire($redisKey, 86400); // 1天
return true; // ACK success
} catch (\Exception $e) {
Redis::del($redisKey);
throw $e; // 重试
}
}
}
| 方案 | 可靠性 | 性能 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 数据库唯一索引 | 数据落库的业务 | |||
| Redis锁 | 高并发、非数据库 | |||
| 业务状态标记 | 状态机类业务 | |||
| 幂等中间件表 | 多服务通用去重 |
最后提醒:消息幂等性并非100%杜绝所有重复,但结合可靠的生产者重试机制、消费者幂等校验、最终一致性方案,可以做到业务意义上的不重复(如最终订单只支付一次)。