简介

基于 ThinkPHP 6.1 构建,集成 Swoole、RabbitMQ 与 Redis,打造高并发、高可用的异步订单与库存处理系统。

利用 Swoole 实现常驻内存、异步非阻塞的任务处理;
通过 RabbitMQ 解耦核心业务流程,确保订单创建、库存扣减、超时取消等操作的可靠异步执行;
借助 Redis 的原子操作(如 Lua 脚本或 INCR/DECR)实现库存预扣与并发控制,有效防止超卖,保障数据一致性。

流程图

请添加图片描述

订单系统流程

订单创建

  • 系统生成订单,并将订单状态设为“待支付”。
  • 订单延时取消队列
  • 锁定库存

支付处理

  • 支付成功,发送库存扣减消息。
  • 扣减库存

超时取消

  • 更新订单状态
  • 释放锁定的商品库存

消息队列应用(基于RabbitMQ)

  • 延迟队列
    使用于订单超时取消逻辑,可以设定一个特定时间后执行取消操作的消息。
  • 死信队列
    处理支付失败、库存扣减失败等异常情况。这些消息可以被转发到死信队列中进行进一步的人工或自动化处理。

扩展安装与配置

安装swoole 和 rabbitmq 扩展

执行

    composer require topthink/think-swoole
    composer require php-amqplib/php-amqplib

swoole 配置

<?php

return [
    'http'       => [
        'enable'     => true,
        'host'       => '0.0.0.0',
        'port'       => 9501,
        'worker_num' => swoole_cpu_num(),
        'options'    => [],
    ],
    'websocket'  => [
        'enable'        => false,
        'route' => true,
        'handler'       => \think\swoole\websocket\Handler::class,
        'ping_interval' => 25000,
        'ping_timeout'  => 60000,
        'room'          => [
            'type'  => 'table',
            'table' => [
                'room_rows'   => 8192,
                'room_size'   => 2048,
                'client_rows' => 4096,
                'client_size' => 2048,
            ],
            'redis'  => [
        'host'          => 'redis',
        'port'          => 6379,
        'max_active'    => 20,
        'max_wait_time' => 5,
        'min_active'    => 1,
    ],
        ],
        'listen'        => [],
        'subscribe'     => [],
    ],
    'rpc'        => [
        'server' => [
            'enable'     => false,
            'host'       => '0.0.0.0',
            'port'       => 9000,
            'worker_num' => swoole_cpu_num(),
            'services'   => [],
        ],
        'client' => [],
    ],
    //队列
    'queue'      => [
        'enable'  => false,
        'workers' => [],
    ],
    'hot_update' => [
        'enable'  => env('APP_DEBUG', false),
        'name'    => ['*.php'],
        'include' => [app_path()],
        'exclude' => [],
    ],
    //连接池
    'pool'       => [
        'db'    => [
            'enable'        => true,
            'max_active'    => 3,
            'max_wait_time' => 5,
        ],
        'cache' => [
            'enable'        => false,
            'max_active'    => 3,
            'max_wait_time' => 5,
        ],
        'redis' => [
            'enable'        => true,
            'max_active'    => 20,
            'max_wait_time' => 5,
            'min_active'    => 1,
        ],
        //自定义连接池
    ],
    'ipc'        => [
        'type'  => 'unix_socket',
        'redis' => [
            'host'          => 'redis',
            'port'          => 6379,
            'max_active'    => 20,
            'max_wait_time' => 5,
            'min_active'    => 1,
        ],
    ],
    //'lock'       => [
        'enable' => false,
        'type'   => 'table',
        'redis'  => [
            'host'          => 'redis',
            'port'          => 6379,
            'max_active'    => 20,
            'max_wait_time' => 5,
            'min_active'    => 1,
        ],
    ],
    'tables'     => [],
    //每个worker里需要预加载以共用的实例
    'concretes'  => [],
    //重置器
    'resetters'  => [],
    //每次请求前需要清空的实例
    'instances'  => [],
    //每次请求前需要重新执行的服务
    'services'   => [],
];

rabbitmq 配置

<?php
return [
    // RabbitMQ 连接配置
    'host'     => env('RABBITMQ_HOST', 'rabbitmq'),
    'port'     => env('RABBITMQ_PORT', 5672),
    'user'     => env('RABBITMQ_USER', 'tp5'),        // 👈 改成你的用户
    'password' => env('RABBITMQ_PASSWORD', 'tp5_secret'), // 👈 改成你的密码
    'vhost'    => env('RABBITMQ_VHOST', 'tp5'),       // 👈 改成你的 vhost
    'keepalive'=> false,

    // 👇 队列配置 —— 【关键修复】'exchange' 字段直接填写交换机真实名称!
    'queues' => [
        // 订单创建
        'order_created' => [
            'name'        => 'order_created',
            'exchange'    => 'order.events.exchange',     // ✅ 修复:真实交换机名
            'routing_key' => 'order.created',
            'durable'     => true,
            'auto_delete' => false,
            'retry_delay' => 5000,
            'max_retries' => 3,
            'dlx_name'    => 'dlx.exchange',
        ],
        
        // 库存扣减
        'inventory_deduct' => [
            'name'        => 'inventory_deduct',
            'exchange'    => 'inventory.events.exchange', // ✅ 修复:真实交换机名
            'routing_key' => 'inventory.deduct',
            'durable'     => true,
            'auto_delete' => false,
            'retry_delay' => 10000,
            'max_retries' => 5,
            'dlx_name'    => 'dlx.exchange',
        ],
        
        // 订单超时(延迟队列)
        'order_timeout' => [
            'name'        => 'order_timeout',
            'exchange'    => 'order.timeout.exchange',    // ✅ 修复:真实交换机名
            'routing_key' => 'order.timeout',
            'durable'     => true,
            'auto_delete' => false,
            'retry_delay' => 15000,
            'max_retries' => 2,
            'dlx_name'    => 'dlx.exchange',
        ],
        
        // 库存回滚
        'inventory_rollback' => [
            'name'        => 'inventory_rollback',
            'exchange'    => 'inventory.events.exchange', // ✅ 修复:真实交换机名
            'routing_key' => 'inventory.rollback',
            'durable'     => true,
            'auto_delete' => false,
            'retry_delay' => 8000,
            'max_retries' => 3,
            'dlx_name'    => 'dlx.exchange',
        ],
        
        // 支付处理
        'payment_processed' => [
            'name'        => 'payment_processed',
            'exchange'    => 'main.exchange',             // ✅ 修复:真实交换机名
            'routing_key' => 'payment.processed',
            'durable'     => true,
            'auto_delete' => false,
            'retry_delay' => 5000,
            'max_retries' => 3,
            'dlx_name'    => 'dlx.exchange',
        ],
        
        // 👇 统一死信队列
        'global.dlq' => [
            'name'        => 'global.dlq',
            'exchange'    => 'dlx.exchange',              // ✅ 修复:真实交换机名
            'routing_key' => '#',
            'durable'     => true,
            'auto_delete' => false,
        ],
    ],

    // 👇 交换机配置(保留用于声明,脚本不再用于绑定映射)
    'exchanges' => [
        'order_events' => [
            'name'        => 'order.events.exchange',
            'type'        => 'topic',
            'durable'     => true,
            'auto_delete' => false,
        ],
        
        'inventory_events' => [
            'name'        => 'inventory.events.exchange',
            'type'        => 'topic',
            'durable'     => true,
            'auto_delete' => false,
        ],
        
        'order_delayed_events' => [
            'name'        => 'order.timeout.exchange',
            'type'        => 'x-delayed-message',
            'durable'     => true,
            'auto_delete' => false,
            'arguments'   => ['x-delayed-type' => 'topic'],
        ],
        
        'dlx' => [
            'name'        => 'dlx.exchange',
            'type'        => 'topic',
            'durable'     => true,
            'auto_delete' => false,
        ],
        
        'main' => [
            'name'        => 'main.exchange',
            'type'        => 'topic',
            'durable'     => true,
            'auto_delete' => false,
        ],
    ],

    // 👇 死信消费者配置(已兼容真实名称)
    'dlx_consumer' => [
        'queue'       => 'global.dlq',    // 👈 和上面队列名一致
        'exchange'    => 'dlx.exchange',  // 👈 和 dlx 交换机名一致
        'routing_key' => '#',
    ],
     'exchange_names' => [
        'order_events'         => 'order.events.exchange',
        'inventory_events'     => 'inventory.events.exchange',
        'order_delayed_events' => 'order.timeout.exchange',
        'main'                 => 'main.exchange',
        'dlx'                  => 'dlx.exchange',
    ],
];

启动swoole

启动:php think swoole
重启:php think swoole restart
停止:php think swoole stop
状态:php think swoole status

启动rabbitmq

php think rabbitmq:consume

注册命令

在config中console.php

<?php
// +----------------------------------------------------------------------
// | 控制台配置
// +----------------------------------------------------------------------
return [
    // 指令定义
    'commands' => [
        'inventory:consumer' => 'app\command\InventoryConsumer',
        'inventory:consumer-simple' => 'app\command\InventoryConsumerSimple',
        'order:consumer' => 'app\command\OrderConsumer',
        'order:timeout-consumer' => 'app\command\OrderTimeoutConsumer',
        'test:timeout-order' => 'app\command\TestTimeoutOrder',
        'dlx:consumer' => 'app\command\DlxConsumer',
        'rabbitmq:health' => 'app\command\RabbitMQHealthCheck',
    ],
];

下单流程

创建订单

// 核心事务:创建订单 + 写入消息
    protected function createOrderInTransaction($userId, $addressId, $remark, $items, $inventoryDeducts, $mode)
    {
        return Db::transaction(function () use ($userId, $addressId, $remark, $items, $inventoryDeducts, $mode) {
            $totalAmount = 0;
            $orderItems = [];
            $messages = [];

            // 准备订单项和消息
            foreach ($items as $item) {
                $sku = $item['sku_model'] ?? GoodsSku::find($item['sku_id']);
                $quantity = $item['quantity'] ?? $item['quantity'];

                if (!$sku) throw new \Exception("商品不存在");

                $total = $sku["price"] * $quantity;
                $totalAmount += $total;

                $orderItems[] = [
                    'goods_id' => $sku["goods_id"],
                    'sku_id' => $sku["id"],
                    'goods_name' => $sku["goods"]["name"] ?? '商品',
                    'sku_specs' => json_encode($sku["specs"] ?? [], JSON_UNESCAPED_UNICODE),
                    'price' => $sku["price"],
                    'quantity' => $quantity,
                    'total_price' => $total,
                ];
                Cache::set("goods_sku:{$sku["id"]}", [
                    'id' => $sku["id"],
                    'goods_id' => $sku["goods_id"],
                    'title' => $sku["title"]??"商品",
                    'price' => $sku["price"],
                ], 300);

                //  准备库存扣减消息
                if (in_array($mode, ['local_message', 'dual'])) {
                    $messages[] = [
                        'message_id' => 'inv_' . uniqid() . '_' . $sku["id"],
                        'exchange' => 'inventory',
                        'routing_key' => 'deduct',
                        'body' => json_encode([
                            'order_id' => 0,
                            'sku_id' => $sku["id"],
                            'quantity' => $quantity,
                        ], JSON_UNESCAPED_UNICODE),
                        'status' => 0,
                        'try_count' => 0,
                        'next_retry_time' => date('Y-m-d H:i:s'),
                    ];
                }
            }
          
            // 创建订单
            $orderNo = Order::generateOrderNo();
            $order = new Order();
            $order->order_no = $orderNo;
            $order->user_id = $userId;
            $order->address_id = $addressId;
            $order->total_amount = $totalAmount;
            $order->pay_amount = $totalAmount;
            $order->status = Order::STATUS_PENDING;
            $order->pay_status = Order::PAY_STATUS_UNPAID;
            $order->remark = $remark;
            $order->save();

            $this->log('info', "创建订单成功,order_id: {$order->id}, mode: {$mode}");
            // 填充订单ID
            foreach ($messages as &$msg) {
                $msg['body'] = str_replace('"order_id":0', '"order_id":' . $order->id, $msg['body']);
            }

            // 批量插入订单项
            foreach ($orderItems as &$item) {
                $item['order_id'] = $order->id;
                $inventoryDeducts[] = array(
                    'order_id' => $order->id,
                    'sku_id' => $item['sku_id'],
                    'quantity' => $item['quantity']
                );
            }
            // 使用insertAll方法进行高效的批量插入
            $insertResult = OrderItem::insertAll($orderItems);
            if (!$insertResult) {
                throw new \Exception('订单项批量插入失败');
            }
            $this->log('info', "插入商品成功,order_id: {$order->id}, insertResult: {$insertResult}");

            

            // 写入本地消息表(local_message 或 dual 模式)
            if (in_array($mode, ['local_message', 'dual'])) {
                foreach ($messages as $msg) {
                    Db::table('local_message')->insert($msg);
                }

                Db::table('local_message')->insert([
                    'message_id' => 'order_created_' . $order->id,
                    'exchange' => 'order',
                    'routing_key' => 'created',
                    'body' => json_encode(['order_id' => $order->id], JSON_UNESCAPED_UNICODE),
                    'status' => 0,
                    'try_count' => 0,
                    'next_retry_time' => date('Y-m-d H:i:s'),
                ]);

                Db::table('local_message')->insert([
                    'message_id' => 'order_timeout_' . $order->id,
                    'exchange' => 'order',
                    'routing_key' => 'timeout',
                    'body' => json_encode(['order_id' => $order->id, 'delay_minutes' => 1], JSON_UNESCAPED_UNICODE),
                    'status' => 0,
                    'try_count' => 0,
                    'next_retry_time' => date('Y-m-d H:i:s', time() + 60),
                ]);
            }
              $this->log('info', "创建订单开始,mode: {$mode}");

            //  同步发送 RabbitMQ(仅在 redis / dual 模式,且不失败时不阻塞)
            // 安全发送消息
            if (in_array($mode, ['redis', 'dual'])) {
                 $this->log('info', "开始处理库存消息");
                $producer = new MessageProducerService();
                $this->log('info', "new MessageProducerService");
                $mqSuccess = true;
                $mqSuccess &= $this->safePublish($producer, 'publishOrderCreated', $order->id);
               $this->log('info', '订单创建,库存扣减遍历 inventoryDeducts='.count($inventoryDeducts));
                foreach ($inventoryDeducts as $item) {
                    $this->log('info', '订单创建,库存扣减消息发送>'.$item["sku_id"].' quantity>'.$item["quantity"]);

                     $stock = $this->redis->get("stock:sku:{$item['sku_id']}");
                      $this->log('send info', '商品 key>' . "stock:sku:{$item['sku_id']}");
                    $this->log('send info', '商品 stock>' . json_encode($stock));

                    $mqSuccess &= $this->safePublish(
                        $producer,
                        'publishInventoryDeduct',
                        $order->id,
                        $item['sku_id'],
                        $item['quantity']
                    );
                }

                 $mqSuccess &= $this->safePublish($producer, 'publishOrderTimeout', $order->id, 60);

                // 如果是纯 Redis 模式且 MQ 失败,应触发回滚
                if ($mode === 'redis' && !$mqSuccess) {
                    throw new \Exception('消息队列异常,订单已回滚');
                }
                
                // dual 模式下,即使 MQ 失败,本地消息表已兜底
            }

            return [
                'success' => true,
                'order_id' => $order->id,
                'order_no' => $orderNo,
                'total_amount' => $totalAmount,
                'item_count' => count($items),
                'mode' => $mode,
            ];
        });
    }

预扣库存

 // 核心:Redis 预扣库存(原子 Lua)
    protected function validateAndDeductRedisStock($items)
    {
        $inventoryDeducts = [];
        // 提取所有 sku_id 并过滤掉可能的空值,然后批量查询商品 SKU 信息并按 ID 索引
        $skuIds = array_filter(array_column($items, 'sku_id'));
       
        if (!empty($skuIds)) {
            // $skus = GoodsSku::whereIn('id', $skuIds)->get()->keyBy('id');
            $skus = GoodsSku::whereIn('id', $skuIds)->select()->toArray();
            $skus = collect(array_column($skus, null, 'id'));

        } else {
            $skus = collect([]);
        }
      
        foreach ($items as $item) {
            $skuId = $item['sku_id'];
            $quantity = $item['quantity'];
            $sku = $skus[$skuId] ?? null;
            if (!$sku) throw new \Exception("商品不存在: {$skuId}");
            if ($quantity <= 0) throw new \Exception("购买数量错误");

            $redisKey = "stock:sku:{$skuId}";
            $luaScript = <<<LUA
if redis.call('GET', KEYS[1]) == false then
    redis.call('SET', KEYS[1], ARGV[1])
end
local stock = tonumber(redis.call('GET', KEYS[1]))
local deduct = tonumber(ARGV[2])
if stock >= deduct then
    redis.call('DECRBY', KEYS[1], deduct)
    return 1
else
    return 0
end
LUA;
            $result = $this->redis->eval($luaScript, [$redisKey, $sku["stock"], $quantity],1);
            
              $stock = $this->redis->get("stock:sku:{$skuId}");

              $this->log('info', '商品 stock>' . json_encode($stock));
             

            if ($result != 1) {
                throw new \Exception("商品 {$sku["id"]} 库存不足");
            }

            $inventoryDeducts[] = [
                'sku_id' => $skuId,
                'quantity' => $quantity,
                'sku_model' => $sku,
                'redis_key' => $redisKey,
            ];

             $this->log('info', '商品 ' . $sku["id"] .' redis_key:'.$redisKey. ' 库存预扣, 剩余库存: ' . $this->redis->get($redisKey));
            // 设置过期时间,防止长期占用
            $this->redis->expire($redisKey, Config::get('order.redis_stock_ttl', 86400));
        }
        return $inventoryDeducts;
    }

订单支付

// 支付订单 - 使用 think-swoole 协程优化
    public function payOrder(string $orderNo, string $paymentMethod = 'wechat'): array
    {
        return Db::transaction(function () use ($orderNo, $paymentMethod) {
            try {
                // 使用缓存优化订单查询
                $cacheKey = "order:{$orderNo}";
                $orderData = Cache::get($cacheKey);
                
                if (!$orderData) {
                    $order = Order::where('order_no', $orderNo)->lock(true)->find();
                    if (!$order) {
                        throw new \Exception('订单不存在');
                    }
                    $orderData = $order->toArray();
                    Cache::set($cacheKey, $orderData, 300);
                }
                
                $order = new Order($orderData);
                
                if (!$order->canPay()) {
                    throw new \Exception('订单不可支付');
                }
                
                // 确保订单状态为待支付
                if ($order->status != \app\model\Order::STATUS_PENDING) {
                    throw new \Exception('订单状态异常');
                }

                // 创建支付记录
                $paymentRecord = new PaymentRecord();
                $paymentRecord->order_id = $order->id;
                $paymentRecord->payment_no = $this->generatePaymentNo();
                $paymentRecord->payment_method = $paymentMethod;
                $paymentRecord->amount = $order->pay_amount;
                $paymentRecord->status = 0; // 待支付
                $paymentRecord->save();

                $paymentResult =$this->processPayment($paymentRecord);

                if ($paymentResult['success']) {
                    // 更新支付记录
                    $paymentRecord->status = PaymentRecord::STATUS_SUCCESS;
                    $paymentRecord->transaction_id = $paymentResult['transaction_id'];
                    $paymentRecord->paid_at = date('Y-m-d H:i:s');
                    $paymentRecord->save();

                    // 更新订单状态
                    $updateResult = Order::where('order_no', $orderNo)
                        ->where('pay_status', Order::PAY_STATUS_UNPAID)
                        ->update([
                            'pay_status' => Order::PAY_STATUS_PAID,
                            'status' => Order::STATUS_PAID,
                            'pay_time' => date('Y-m-d H:i:s')
                        ]);

                    if (!$updateResult) {
                        throw new \Exception('订单状态更新失败');
                    }

                    // 清除缓存
                    Cache::delete($cacheKey);

                    // 并发处理后续任务
                    $tasks = [];
                    
                    // 异步任务1:处理支付后业务
                    $tasks[] = Coroutine::create(function () use ($order) {
                        try {
                            $this->afterPayment($order);
                        } catch (\Exception $e) {
                            error_log('[ERROR] 支付后处理失败: ' . $e->getMessage());
                        }
                    });

                    // 异步任务2:发送支付成功通知
                    $tasks[] = Coroutine::create(function () use ($order) {
                        try {
                            $this->sendPaymentSuccessNotification($order);
                        } catch (\Exception $e) {
                            error_log('[ERROR] 支付通知发送失败: ' . $e->getMessage());
                        }
                    });

                    // 异步任务3:记录支付日志
                    $tasks[] = Coroutine::create(function () use ($paymentRecord) {
                        try {
                            error_log('[INFO] 支付成功记录: ' . json_encode([
                                'payment_id' => $paymentRecord->id,
                                'amount' => $paymentRecord->amount,
                                'method' => $paymentRecord->payment_method
                            ]));
                        } catch (\Exception $e) {
                            error_log('[ERROR] 支付日志记录失败: ' . $e->getMessage());
                        }
                    });

                    return [
                        'success' => true,
                        'order_no' => $order->order_no,
                        'payment_no' => $paymentRecord->payment_no,
                        'transaction_id' => $paymentResult['transaction_id'],
                        'paid_amount' => $order->pay_amount
                    ];
                } else {
                    // 支付失败
                    $paymentRecord->status = PaymentRecord::STATUS_FAILED;
                    $paymentRecord->error_message = $paymentResult['error_message'] ?? '支付失败';
                    $paymentRecord->save();

                    throw new \Exception($paymentResult['error_message'] ?? '支付失败');
                }

            } catch (\Exception $e) {
                $this->log('error', '订单支付失败: ' . $e->getMessage());
                throw $e;
            }
        });
    }

库存扣减

 /**
     * 处理库存扣减
     */
    public function handleInventoryDeduct(array $data): bool
    {
        $orderId  = $data['order_id'] ?? 0;
        $skuId    = $data['sku_id'] ?? 0;
        $quantity = $data['quantity'] ?? 0;
        $this->safeLog('info', "Redis 获取key》"."stock:sku:{$skuId}");
        // 从 Redis 获取库存
        $stock = $this->redis->get("stock:sku:{$skuId}");
        $this->safeLog('info', "Redis 中库存:".$stock);
        if (!$stock) {
            $this->safeLog('error', "Redis 中未找到库存 kkuid>".$skuId);
            return false;
        }

        if (!$orderId || !$skuId || !$quantity) {
            $this->safeLog('error', '库存扣减参数错误', ['data' => $data]);
            return false;
        }

        // 分布式锁
        $lockKey = "inventory_deduct_lock:{$skuId}";
        if (!$this->redis->set($lockKey, 1, ['nx', 'ex' => 5])) {
            $this->safeLog('warning', "库存扣减锁定失败", ['sku_id' => $skuId]);
            return false;
        }

        try {
            $sku = GoodsSku::find($skuId);
            if (!$sku) {
                $this->safeLog('error', "商品不存在", ['sku_id' => $skuId]);
                return false;
            }

            if ($sku->stock < $quantity) {
                $this->safeLog('error', "库存不足", [
                    'sku_id' => $skuId,
                    'need'   => $quantity,
                    'stock'  => $sku->stock,
                ]);
                return false;
            }

            // 乐观锁更新库存
            $updated = GoodsSku::where('id', $skuId)
                ->where('stock', '>=', $quantity)
                ->dec('stock', $quantity)
                ->update();

            if ($updated) {
                $this->safeLog('info', "库存扣减成功", [
                    'sku_id'   => $skuId,
                    'quantity' => $quantity,
                    'order_id' => $orderId,
                ]);
                Cache::delete("goods_sku:{$skuId}");
                return true;
            }

            $this->safeLog('error', "库存扣减失败", ['sku_id' => $skuId]);
            return false;

        } catch (\Throwable $e) {
            $this->safeLog('error', "库存扣减异常", [
                'error' => $e->getMessage(),
                'trace' => $e->getTraceAsString(),
                'data'  => $data,
            ]);
            return false;

        } finally {
            $this->redis->del($lockKey);
        }
    }

超时订单取消

/**
     * 处理订单超时消息
     */
    
    public function handleOrderTimeout($data = null): bool
    {
       
        $orderId = $data['order_id'] ?? 0;
        
        if (!$orderId) {
            Log::error('订单超时消息参数错误', $data);
            return false;
        }
        
        $result = Db::transaction(function () use ($orderId) {
            $order = Order::where('id', $orderId)
                ->where('status', Order::STATUS_PENDING)
                ->where('pay_status', Order::PAY_STATUS_UNPAID)
                ->lock(true)
                ->find();
            
            if (!$order) {
                Log::info("订单已支付或不存在,无需超时处理", ['order_id' => $orderId]);
                return false;
            }
            
            // 更新订单状态为已取消
            $updated = Order::where('id', $orderId)
                ->update([
                    'status' => Order::STATUS_CANCELLED,
                    'cancel_reason' => '订单超时未支付',
                    'cancelled_at' => date('Y-m-d H:i:s')
                ]);
            
            if ($updated) {
                Log::info("订单超时取消成功", ['order_id' => $orderId]);
                
                // 获取订单商品信息
                $orderItems = $order->items()->select();
                
                // 发送库存回滚消息
                $producer = new MessageProducerService();
                foreach ($orderItems as $item) {
                    $producer->publishInventoryRollback(
                        $orderId,
                        $item['sku_id'],
                        $item['quantity']
                    );
                }
                
                return true;
            }
            
            return false;
        });
        
        if ($result) {
            Log::info("订单超时处理完成", ['order_id' => $orderId]);
        }
        
        return $result;
    }

库存回滚

/**
     * 处理库存回滚
     */
    public function handleInventoryRollback(array $data): bool
    {
        // 记录收到的消息完整结构
        $this->safeLog('info', '收到库存回滚消息', ['data_structure' => array_keys($data)]);
        $this->safeLog('info', '库存回滚消息完整内容', ['data' => $data]);
        
        $orderId  = $data['order_id'] ?? 0;
        
        // 兼容两种数据格式:直接参数或items数组
        if (isset($data['items']) && is_array($data['items']) && !empty($data['items'])) {
            $this->safeLog('info', '使用items数组格式处理库存回滚');
            // 处理items数组格式
            $item = $data['items'][0]; // 假设一次只处理一个商品
            $skuId = $item['sku_id'] ?? 0;
            $quantity = $item['quantity'] ?? 0;
        } else {
            $this->safeLog('info', '使用直接参数格式处理库存回滚');
            // 处理直接参数格式
            $skuId = $data['sku_id'] ?? 0;
            $quantity = $data['quantity'] ?? 0;
        }

        $this->safeLog('info', '解析后的库存回滚参数', ['order_id' => $orderId, 'sku_id' => $skuId, 'quantity' => $quantity]);

        if (!$orderId || !$skuId || !$quantity) {
            $this->safeLog('error', '库存回滚参数错误', ['data' => $data]);
            return false;
        }

        $lockKey = "inventory_rollback_lock:{$skuId}";
        if (!$this->redis->set($lockKey, 1, ['nx', 'ex' => 5])) {
            $this->safeLog('warning', "库存回滚锁定失败", ['sku_id' => $skuId]);
            return false;
        }

        try {
            $updated = GoodsSku::where('id', $skuId)
                ->inc('stock', $quantity)
                ->update();

            if ($updated) {
                $this->safeLog('info', "库存回滚成功", [
                    'sku_id'   => $skuId,
                    'quantity' => $quantity,
                    'order_id' => $orderId,
                ]);
                Cache::delete("goods_sku:{$skuId}");
                return true;
            }

            $this->safeLog('error', "库存回滚失败", ['sku_id' => $skuId]);
            return false;

        } catch (\Throwable $e) {
            $this->safeLog('error', "库存回滚异常", [
                'error' => $e->getMessage(),
                'trace' => $e->getTraceAsString(),
                'data'  => $data,
            ]);
            return false;

        } finally {
            $this->redis->del($lockKey);
        }
    }

查看完整demo

Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐