'orderChange', StarChainMessageEnum::ORDER_DELIVERY => 'orderDelivery', StarChainMessageEnum::REFUND_CHANGE => 'refundOrderChange', StarChainMessageEnum::GOODS_CHANGE => 'goodsChange', StarChainMessageEnum::GOODS_REMOVE => 'goodsPulled', StarChainMessageEnum::ATTR_REMOVE => 'goodsChange', StarChainMessageEnum::ATTR_CHANGE => 'goodsChange', StarChainMessageEnum::SELECTION_REMOVE => 'goodsPulled', StarChainMessageEnum::SELECTION_GOODS_ADD => 'notHandler', StarChainMessageEnum::ADD_SELECTION_GOODS => 'notHandler', StarChainMessageEnum::REMOVE_SELECTION_GOODS => 'notHandler', StarChainMessageEnum::SELECTION_PRICE_CHANGE => 'goodsChange', StarChainMessageEnum::GOODS_UP_AND_DOWN => 'goodsChange', StarChainMessageEnum::ORDER_CREATE => 'notHandler', StarChainMessageEnum::ORDER_CHANGE_DELIVERY => 'notHandler', ]; /** * 处理成功的消息IDs * [ * id, * id, * ] * * @var array */ protected $successIds = []; /** * 处理失败的消息 * [ * id => msg, * id => msg, * ] * * @var array */ protected $failMsgs = []; /** * 使用推送消息进行发货 * * @return string */ protected $msgDeliveryType = true; /** * 获取一个消息rediskey * * @return string */ private function getMessageRedisKey() :string { return 'star-chain-message-list:' . $this->mall_id; } /** * 获取消息池消息加入redis队列 * * @return void */ public function getMessagePool() { // 未安装 if (empty($this->api)) return false; // 全部消息 try { $msg_list = $this->api->messagePool(); } catch (Exception $e) { Yii::error( "获取消息失败,mall" . $this->mall_id . ", 失败内容为:" . $e->getMessage() ); } if (empty($msg_list)) return false; // $msg_list = array_slice($msg_list, 0, 50); try { // 加入队列 $redis = Yii::$app->redis; $msg_ids = []; foreach ($msg_list as $v) { $redis->rpush($this->getMessageRedisKey(), Json::encode($v)); $msg_ids[] = $v['id']; } $this->api->removeMessagePoolByParam($msg_ids); } catch (Exception $e) { Yii::error( "删除消息失败ID为" . Json::encode($msg_ids) . ", 失败内容为:" . $e->getMessage() ); } } /** * 获取消息插入数据库 * * @return void */ public function getMessage() { Yii::$app->services->rabbitMq->push( RabbitMqEnum::EXCHANGE_CHAIN, RabbitMqEnum::CHAIN_QUEUE, ['mall_id' => $this->mall_id], self::class, 'handleSyncMsg' ); } /** * 异步执行 * * @return void */ public static function handleSyncMsg($data) { $service = new static(['mall_id' => $data['mall_id']]); $service->apiLoadMsgToDB(); } /** * api获取消息并存储 * * @return void */ public function apiLoadMsgToDB() { // 未安装 if (empty($this->api)) return false; // 全部消息 try { $msg_list = $this->api->messagePool(); } catch (Exception $e) { Yii::error( "获取消息失败,mall" . $this->mall_id . ", 失败内容为:" . $e->getMessage() ); } if (empty($msg_list)) return false; // 消息ID $msg_ids = []; try { // 写入数据库 Message::addMsg($this->mall_id, $msg_list); // 加入队列 foreach ($msg_list as $item) { $item['mall_id'] = $this->mall_id; // 延迟执行 Yii::$app->services->rabbitMq->delayTask( static::DELAY, $item, self::class, 'handleSync', 'chain' ); $msg_ids[] = $item['id']; } // 删除消息 $this->api->removeMessagePoolByParam($msg_ids); } catch (Exception $e) { Yii::error('星链供应链消息存储失败'); Yii::error($msg_list); } } /** * 获取消息池消息(redis消息队列) * * @return void */ public function messagePoolHandler() { // 获取消息 $redis = Yii::$app->redis; $msg_list = $redis->eval($this->_lua_get_msg(), 2, $this->getMessageRedisKey(), self::LIMIT); if (empty($msg_list)) return false; $msg_list = array_map(function ($item) { return Json::decode($item); }, $msg_list); // 储存消息 Message::addMsg($this->mall_id, $msg_list); $msg_ids = array_column($msg_list, 'id'); $msg_list = $this->groupMsg($msg_list); $this->messageListHandle($msg_list); } /** * 消息处理 * * @param array $data * @return void */ public static function handleSync(array $data) { $service = new static(['mall_id' => $data['mall_id']]); $msg_list = $service->groupMsg([$data]); $service->messageListHandle($msg_list); } /** * 处理消息 * * @param array $msg_list * @return void */ public function messageListHandle($msg_list) { // $db = Yii::$app->db->beginTransaction(); try { foreach ($msg_list as $enum => $v) { // echo "消息类型{$dasprid-enum}, 数量" . count($v); // 按类型进行消息处理 if ($func = $this->ex_func[$enum]) { // 处理消息 $this->$func($v); } } // 完成处理 $this->finish(); // $db->commit(); } catch (Exception $e) { Yii::error('星链供应链消息处理失败: ' . json_encode([ 'file' => $e->getFile(), 'msg' => $e->getMessage(), 'line' => $e->getLine(), ])); Yii::error($msg_list); // $db->rollBack(); } } /** * 完成消息处理 * * @return void */ public function finish() { // 处理成功 Message::setHandled($this->successIds); // 处理失败 foreach ($this->failMsgs as $fail_msg) { Message::setHandledFail($fail_msg['ids'], $fail_msg['msg']); } $this->clearMsgs(); } /** * 消息按内容分组 * * @param array $msgs 消息数组 * @return array */ public function groupMsg($msgs) :array { $group_msg = []; foreach ($msgs as $v) { $group_msg[$v['enumNum']][] = $v; } return $group_msg; } /** * 选品移除 * * @param array $msg_list * @return void */ public function selectionRemove($msg_list) { $service = new SpuService(['mall_id' => $this->mall_id]); foreach ($msg_list as $v) { // 更新的SPU_ID $content = Json::decode($v['content']); $spu_ids[] = strval($content['spuId']); } $service->removeImportSpuIds($spu_ids); } /** * 添加选品 * * @param array $msg_list * @return void */ public function selectionAdd($msg_list) { // 添加选品不做操作 } /** * 商品下架 * * @param array $msg_list * @return void */ public function goodsPulled($msg_list) { $ids = []; try { $spu_ids = []; foreach ($msg_list as $v) { // 更新的SPU_ID $content = Json::decode($v['content']); $spu_ids[] = strval($content['spuId']); $ids[] = $v['id']; } // 更新商城商品数据 $service = new GoodsService(['mall_id' => $this->mall_id]); $service->goodsPulled($spu_ids); } catch (Exception $e) { return $this->addFailMsgs($ids, $e->getMessage()); } return $this->addSuccessIds($ids); } /** * 订单发货 * * @param array $msg_list * @return array */ public function orderDelivery($msg_list) { // 使用api查询方式 if (!$this->msgDeliveryType) { $ids = []; $order_sn = []; foreach ($msg_list as $v) { // 更新的SPU_ID $content = Json::decode($v['content']); $order_sn[] = strval($content['orderSn']); $ids[] = $v['id']; } try { $service = new OrderService(['mall_id' => $this->mall_id]); $service->orderDeliveryByOrderSn($order_sn); } catch (Exception $e) { return $this->addFailMsgs($ids, $e->getMessage()); } return $this->addSuccessIds($ids); } // 直接使用推送数据 if ($this->msgDeliveryType) { $service = new OrderService(['mall_id' => $this->mall_id]); foreach ($msg_list as $v) { // 推送订单数据 $content = Json::decode($v['content']); $order_sn = strval($content['orderSn']); $delivery_sn = strval($content['deliverySn']); $sku_ids = $content['skuIdList']; $delivery_corp_sn = $content['deliveryCorpSn']; try { // 发货 $service->orderDeliveryByPushMsg( $order_sn, $delivery_sn, $sku_ids, $delivery_corp_sn ); } catch (Exception $e) { return $this->addFailMsgs([$v['id']], $e->getMessage()); } return $this->addSuccessIds([$v['id']]); } } } /** * 售后订单状态变更 * * @param array $msg_list * status 售后状态( * 1,待卖家审核 * 2,卖家拒绝退款, * 3,退款成功, * 4,卖家拒绝退货退款, * 5,待买家退货, * 6,买家已退货,待卖家收货, * 7,买家已退货,卖家拒绝收货, * 8,卖家已收货,待确认退款, * 9,退货退款成功, * 10,卖家拒绝换货, * 11,买家已退货,待卖家换货, * 12,卖家已换货,待买家收货, * 13,换货成功, * 14,已关闭, * 15,卖家同意退款, * 16,卖家同意仅退款, * 17,卖家拒绝仅退款, * 18,仅退款成功 * ) * @return void */ public function refundOrderChange($msg_list) { /** * @var RefundService */ $service = new RefundService(['mall_id' => $this->mall_id]); foreach ($msg_list as $v) { // 消息内容 $content = Json::decode($v['content']); // 修改售后信息 if ($service->statusChange($content['returnSn'], $content['status'])) { $this->addSuccessIds([$v['id']]); } else { $this->addFailMsgs([$v['id']], $service->getErrorMsg()); } } } /** * 订单更变事件 * * @param array $msg_list * @return void */ public function orderChange($msg_list) { $service = new OrderService(['mall_id' => $this->mall_id]); foreach ($msg_list as $v) { // 消息内容 $content = Json::decode($v['content']); // 修改订单信息 if ($service->orderStatusChange($content['orderSn'], $content['status'])) { $this->addSuccessIds([$v['id']]); } else { $this->addFailMsgs([$v['id']], $service->getErrorMsg()); } } } /** * 商品信息变更处理 * * @param array $msg_list * @return void */ public function goodsChange($msg_list) { $ids = []; try { $spu_ids = []; foreach ($msg_list as $v) { // 更新的SPU_ID $content = Json::decode($v['content']); $spu_ids[] = strval($content['spuId']); $ids[] = $v['id']; } // 接口获取商品数据 $this->api->setNotExcption(); $spu_list = array_map(function ($item) { return $item[0]; }, $this->api->getSpuBySpuIdsPool($spu_ids)); // 更新供应链商品数据 Goods::saveList($spu_list); // 更新商城商品数据 $service = new GoodsService(['mall_id' => $this->mall_id]); foreach ($spu_list as $spu) { $service->editGoods($spu['spuId']); } } catch (Exception $e) { return $this->addFailMsgs($ids, $e->getMessage()); } return $this->addSuccessIds($ids); } /** * 不做处理 * * @param array $msg_list * @return array */ public function notHandler($msg_list) { $msg_ids = array_column($msg_list, 'id'); return $msg_ids; } /** * lua脚本 * _key redis key * _limit 循环取出数量 * * @return string */ public function _lua_get_msg() { return <<