| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653 |
- <?php
- namespace addons\StarChain\common\services;
- use addons\StarChain\common\enums\StarChainMessageEnum;
- use addons\StarChain\common\models\Goods;
- use addons\StarChain\common\models\Message;
- use common\enums\RabbitMqEnum;
- use Exception;
- use Yii;
- use yii\helpers\Json;
- /**
- * 消息处理服务
- */
- class MessageHandlerService extends BaseService
- {
- /**
- * 每次处理量
- */
- const LIMIT = 20;
- /**
- * 延迟执行时间
- */
- const DELAY = 10;
- /**
- * 处理方法
- *
- * @var array
- */
- protected $ex_func = [
- StarChainMessageEnum::ORDER_CHANGE => '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 <<<SCRIPT
- -- key
- local _key = KEYS[1]
- local _limit = tonumber(KEYS[2])
- local _count = redis.call('LLEN', _key)
- if _limit > _count then
- _limit = _count
- end
- if _count > 0 then
- local msg = {}
- -- 循环_limit
- for i = 1, _limit
- do
- msg[i] = redis.call('LPOP', _key)
- end
- return msg
- else
- return 0
- end
- SCRIPT;
- }
- /**
- * 添加成功msg_id
- *
- * @param array $ids
- * @return void
- */
- public function addSuccessIds(array $ids)
- {
- $this->successIds = array_merge($this->successIds, $ids);
- }
- /**
- * 添加处理失败消息
- *
- * @param array $ids
- * @param string $msg
- * @return void
- */
- public function addFailMsgs(array $ids, string $msg)
- {
- $this->failMsgs[] = [
- 'ids' => $ids,
- 'msg' => $msg
- ];
- }
- /**
- * 清空消息
- *
- * @return void
- */
- public function clearMsgs()
- {
- $this->successIds = [];
- $this->failMsgs = [];
- }
- /**
- * 测试数据
- *
- * @return array
- */
- public function testMsgList()
- {
- $msg_list[] = [
- 'id' => '66346192',
- 'tenantId' => '1535098852173492226',
- 'enumNum' => 8,
- 'content' => '{"spuId":1426102374728585218}',
- 'createTime' => '2022-07-22 14:28:59',
- ];
- // $msg_list = array_map(function ($item) {
- // return Json::decode($item);
- // }, $msg_list);
- $msg_list = $this->groupMsg($msg_list);
- $this->messageListHandle($msg_list);
- /* $msg_list[] = [
- 'id' => '63503428',
- 'tenantId' => '1535098852173492226',
- 'enumNum' => 3,
- 'content' => '{"origStatus":1,"origStatusName":"待卖家审核","returnSn":"SH8763722084","status":15,"statusName":"卖家同意退款"}',
- 'createTime' => '2022-06-27 11:49:06',
- ]; */
- /* $msg_list[] = [
- 'id' => '63609622',
- 'tenantId' => '1535098852173492226',
- 'enumNum' => 1,
- 'content' => '{"orderSn":"PO182297850412888","origStatus":20,"origStatusName":"待发货","status":30,"statusName":"待收货"}',
- 'createTime' => '2022-06-28 10:05:03',
- ];
- $msg_list[] = [
- 'id' => '63609625',
- 'tenantId' => '1535098852173492226',
- 'enumNum' => 2,
- 'content' => '{"orderSn":"PO182297850412888","skuIdList":["1526494051752783874"]}',
- 'createTime' => '2022-06-28 10:05:03',
- ]; */
- /* $msg_list[] = [
- 'id' => '63625249',
- 'tenantId' => '1535098852173492226',
- 'enumNum' => 3,
- 'content' => '{"origStatus":1,"origStatusName":"待卖家审核","returnSn":"SH4873942686","status":5,"statusName":"待买家退货"}',
- 'createTime' => '2022-06-28 11:27:30',
- ]; */
- }
- }
|