本文目录导读:

在 PHP 中避免重复消费(幂等性处理)通常涉及以下几个层面,我会提供从简单到复杂的解决方案,并给出具体的代码示例。
使用唯一标识符(最基础)
基于幂等键(Idempotency Key)
<?php
class PaymentService {
private $redis;
public function __construct($redis) {
$this->redis = $redis;
}
/**
* 处理支付请求
* @param string $orderId 订单ID
* @param float $amount 金额
* @return array
*/
public function processPayment($orderId, $amount) {
// 生成幂等键
$idempotencyKey = "payment:{$orderId}";
// 检查是否已处理(原子操作避免并发问题)
$result = $this->redis->set(
$idempotencyKey,
'processing',
['NX', 'EX' => 300] // NX: 不存在才设置,EX: 5分钟过期
);
if ($result === false) {
// 已经处理过或正在处理
$status = $this->redis->get($idempotencyKey);
return [
'status' => $status === 'processing' ? '重复请求' : '已处理',
'message' => '该订单已处理,请勿重复提交'
];
}
try {
// 执行业务逻辑(支付操作)
$this->executePayment($orderId, $amount);
// 标记为已完成
$this->redis->set($idempotencyKey, 'completed');
return [
'status' => 'success',
'message' => '支付成功'
];
} catch (\Exception $e) {
// 失败时删除标记,允许重试
$this->redis->del($idempotencyKey);
throw $e;
}
}
private function executePayment($orderId, $amount) {
// 实际支付逻辑
}
}
数据库层面的幂等性
使用数据库唯一索引
<?php
class OrderService {
private $pdo;
public function __construct(PDO $pdo) {
$this->pdo = $pdo;
}
/**
* 创建订单(防止重复创建)
*/
public function createOrder($orderId, $userId, $items) {
try {
$sql = "INSERT INTO orders (order_id, user_id, items, created_at)
VALUES (:order_id, :user_id, :items, NOW())";
$stmt = $this->pdo->prepare($sql);
$stmt->execute([
':order_id' => $orderId,
':user_id' => $userId,
':items' => json_encode($items)
]);
return ['status' => 'success', 'order_id' => $orderId];
} catch (PDOException $e) {
// 检查是否因唯一索引冲突
if ($e->getCode() == 23000) { // 唯一约束冲突
return ['status' => 'duplicate', 'message' => '订单已存在'];
}
throw $e;
}
}
}
使用数据库锁
悲观锁实现
<?php
class InventoryService {
private $pdo;
public function __construct(PDO $pdo) {
$this->pdo = $pdo;
}
/**
* 扣减库存(使用悲观锁防止超卖)
*/
public function deductStock($productId, $quantity) {
$this->pdo->beginTransaction();
try {
// 使用 SELECT FOR UPDATE 锁定行
$sql = "SELECT stock FROM products
WHERE id = :product_id
FOR UPDATE";
$stmt = $this->pdo->prepare($sql);
$stmt->execute([':product_id' => $productId]);
$row = $stmt->fetch(PDO::FETCH_ASSOC);
if ($row['stock'] < $quantity) {
$this->pdo->rollBack();
return ['status' => 'failed', 'message' => '库存不足'];
}
// 更新库存
$sql = "UPDATE products
SET stock = stock - :quantity
WHERE id = :product_id";
$stmt = $this->pdo->prepare($sql);
$stmt->execute([
':quantity' => $quantity,
':product_id' => $productId
]);
$this->pdo->commit();
return ['status' => 'success'];
} catch (\Exception $e) {
$this->pdo->rollBack();
throw $e;
}
}
}
乐观锁实现
<?php
class VersionedOrderService {
private $pdo;
public function __construct(PDO $pdo) {
$this->pdo = $pdo;
}
/**
* 更新订单状态(使用版本号防止并发修改)
*/
public function updateOrderStatus($orderId, $newStatus, $version) {
$sql = "UPDATE orders
SET status = :new_status, version = version + 1
WHERE id = :order_id AND version = :version";
$stmt = $this->pdo->prepare($sql);
$result = $stmt->execute([
':new_status' => $newStatus,
':order_id' => $orderId,
':version' => $version
]);
// 如果影响行数为0,说明版本号不匹配(已被其他请求修改)
if ($stmt->rowCount() === 0) {
return ['status' => 'conflict', 'message' => '订单已被其他请求修改'];
}
return ['status' => 'success'];
}
}
完整的生产级方案
<?php
class IdempotentRequestHandler {
private $redis;
private $pdo;
public function __construct($redis, PDO $pdo) {
$this->redis = $redis;
$this->pdo = $pdo;
}
/**
* 处理幂等请求的通用方法
*/
public function handleRequest($requestId, callable $callback, $ttl = 3600) {
// 1. 检查Redis中是否已有处理结果
$cacheKey = "result:{$requestId}";
$cachedResult = $this->redis->get($cacheKey);
if ($cachedResult !== false) {
return json_decode($cachedResult, true);
}
// 2. 获取分布式锁(防止并发处理)
$lockKey = "lock:{$requestId}";
$lockAcquired = $this->redis->set(
$lockKey,
'locked',
['NX', 'EX' => $ttl]
);
if ($lockAcquired === false) {
// 等待或返回重复请求
throw new \Exception('请求正在处理中,请勿重复提交');
}
try {
// 3. 再次检查数据库(双保险)
$existing = $this->checkDatabase($requestId);
if ($existing) {
return $existing;
}
// 4. 执行实际业务逻辑
$result = $callback();
// 5. 将结果存储
$this->saveResult($requestId, $result);
// 6. 缓存结果
$this->redis->set($cacheKey, json_encode($result), 'EX', $ttl);
return $result;
} finally {
// 释放锁
$this->redis->del($lockKey);
}
}
private function checkDatabase($requestId) {
$sql = "SELECT result FROM request_logs WHERE request_id = :request_id";
$stmt = $this->pdo->prepare($sql);
$stmt->execute([':request_id' => $requestId]);
$row = $stmt->fetch(PDO::FETCH_ASSOC);
return $row ? json_decode($row['result'], true) : null;
}
private function saveResult($requestId, $result) {
$sql = "INSERT INTO request_logs (request_id, result, created_at)
VALUES (:request_id, :result, NOW())
ON DUPLICATE KEY UPDATE result = VALUES(result)";
$stmt = $this->pdo->prepare($sql);
$stmt->execute([
':request_id' => $requestId,
':result' => json_encode($result)
]);
}
}
// 使用示例
$handler = new IdempotentRequestHandler($redis, $pdo);
$requestId = $_POST['request_id'] ?? uniqid(); // 客户端生成或服务器生成
$result = $handler->handleRequest($requestId, function() {
// 实际业务逻辑
return processPayment();
});
消息队列中的幂等消费
<?php
class QueueConsumer {
private $redis;
private $pdo;
/**
* 消费RabbitMQ消息(防止重复消费)
*/
public function consumeMessage($messageId, $payload) {
// 1. 使用Redis SetNx确保消息只被处理一次
$processedKey = "processed:{$messageId}";
$isProcessed = $this->redis->set($processedKey, '1', 'NX', 'EX' => 86400); // 24小时
if ($isProcessed === false) {
echo "消息 {$messageId} 已被处理过,跳过\n";
return;
}
// 2. 记录处理日志
$this->logProcessing($messageId, $payload);
try {
// 3. 业务处理
$this->processBusinessLogic($payload);
// 4. 标记成功
$this->redis->set("success:{$messageId}", '1');
} catch (\Exception $e) {
// 失败时删除标记,允许重新消费
$this->redis->del($processedKey);
// 记录失败日志并重新投递
error_log("消息处理失败: {$e->getMessage()}");
$this->requeueMessage($messageId, $payload);
}
}
private function logProcessing($messageId, $payload) {
$sql = "INSERT INTO message_logs (message_id, payload, status, created_at)
VALUES (:message_id, :payload, 'processing', NOW())
ON DUPLICATE KEY UPDATE payload = VALUES(payload)";
$stmt = $this->pdo->prepare($sql);
$stmt->execute([
':message_id' => $messageId,
':payload' => json_encode($payload)
]);
}
private function requeueMessage($messageId, $payload) {
// 重新投递消息
// 根据具体的MQ实现
}
}
- 请求级别:客户端生成唯一
request_id,服务端检查处理 - 数据库级别:使用唯一索引、版本号控制
- 分布式锁:Redis 实现 Set NX EX
- 消息队列:消费前检查处理记录
注意事项
- 原子性:确保检查和处理是原子操作
- 过期时间:设置合理的TTL,防止数据累积
- 重试机制:失败后允许重试,但不允许重复成功
- 日志记录:完整记录处理过程和结果
- 监控告警:及时发现和处理异常情况
选择合适的方案取决于你的具体场景:
- 低并发、简单场景:数据库唯一索引即可
- 高并发、分布式场景:需要 Redis + 数据库组合方案
- 消息队列场景:消费前幂等检查必不可少