PHP 怎么避免重复消费

wen PHP项目 1

本文目录导读:

PHP 怎么避免重复消费

  1. 使用唯一标识符(最基础)
  2. 数据库层面的幂等性
  3. 使用数据库锁
  4. 完整的生产级方案
  5. 消息队列中的幂等消费
  6. 注意事项

在 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实现
    }
}
  1. 请求级别:客户端生成唯一 request_id,服务端检查处理
  2. 数据库级别:使用唯一索引、版本号控制
  3. 分布式锁:Redis 实现 Set NX EX
  4. 消息队列:消费前检查处理记录

注意事项

  • 原子性:确保检查和处理是原子操作
  • 过期时间:设置合理的TTL,防止数据累积
  • 重试机制:失败后允许重试,但不允许重复成功
  • 日志记录:完整记录处理过程和结果
  • 监控告警:及时发现和处理异常情况

选择合适的方案取决于你的具体场景:

  • 低并发、简单场景:数据库唯一索引即可
  • 高并发、分布式场景:需要 Redis + 数据库组合方案
  • 消息队列场景:消费前幂等检查必不可少

抱歉,评论功能暂时关闭!