MessageHandlerService.php 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653
  1. <?php
  2. namespace addons\StarChain\common\services;
  3. use addons\StarChain\common\enums\StarChainMessageEnum;
  4. use addons\StarChain\common\models\Goods;
  5. use addons\StarChain\common\models\Message;
  6. use common\enums\RabbitMqEnum;
  7. use Exception;
  8. use Yii;
  9. use yii\helpers\Json;
  10. /**
  11. * 消息处理服务
  12. */
  13. class MessageHandlerService extends BaseService
  14. {
  15. /**
  16. * 每次处理量
  17. */
  18. const LIMIT = 20;
  19. /**
  20. * 延迟执行时间
  21. */
  22. const DELAY = 10;
  23. /**
  24. * 处理方法
  25. *
  26. * @var array
  27. */
  28. protected $ex_func = [
  29. StarChainMessageEnum::ORDER_CHANGE => 'orderChange',
  30. StarChainMessageEnum::ORDER_DELIVERY => 'orderDelivery',
  31. StarChainMessageEnum::REFUND_CHANGE => 'refundOrderChange',
  32. StarChainMessageEnum::GOODS_CHANGE => 'goodsChange',
  33. StarChainMessageEnum::GOODS_REMOVE => 'goodsPulled',
  34. StarChainMessageEnum::ATTR_REMOVE => 'goodsChange',
  35. StarChainMessageEnum::ATTR_CHANGE => 'goodsChange',
  36. StarChainMessageEnum::SELECTION_REMOVE => 'goodsPulled',
  37. StarChainMessageEnum::SELECTION_GOODS_ADD => 'notHandler',
  38. StarChainMessageEnum::ADD_SELECTION_GOODS => 'notHandler',
  39. StarChainMessageEnum::REMOVE_SELECTION_GOODS => 'notHandler',
  40. StarChainMessageEnum::SELECTION_PRICE_CHANGE => 'goodsChange',
  41. StarChainMessageEnum::GOODS_UP_AND_DOWN => 'goodsChange',
  42. StarChainMessageEnum::ORDER_CREATE => 'notHandler',
  43. StarChainMessageEnum::ORDER_CHANGE_DELIVERY => 'notHandler',
  44. ];
  45. /**
  46. * 处理成功的消息IDs
  47. * [
  48. * id,
  49. * id,
  50. * ]
  51. *
  52. * @var array
  53. */
  54. protected $successIds = [];
  55. /**
  56. * 处理失败的消息
  57. * [
  58. * id => msg,
  59. * id => msg,
  60. * ]
  61. *
  62. * @var array
  63. */
  64. protected $failMsgs = [];
  65. /**
  66. * 使用推送消息进行发货
  67. *
  68. * @return string
  69. */
  70. protected $msgDeliveryType = true;
  71. /**
  72. * 获取一个消息rediskey
  73. *
  74. * @return string
  75. */
  76. private function getMessageRedisKey() :string
  77. {
  78. return 'star-chain-message-list:' . $this->mall_id;
  79. }
  80. /**
  81. * 获取消息池消息加入redis队列
  82. *
  83. * @return void
  84. */
  85. public function getMessagePool()
  86. {
  87. // 未安装
  88. if (empty($this->api)) return false;
  89. // 全部消息
  90. try {
  91. $msg_list = $this->api->messagePool();
  92. } catch (Exception $e) {
  93. Yii::error(
  94. "获取消息失败,mall" . $this->mall_id . ", 失败内容为:" . $e->getMessage()
  95. );
  96. }
  97. if (empty($msg_list)) return false;
  98. // $msg_list = array_slice($msg_list, 0, 50);
  99. try {
  100. // 加入队列
  101. $redis = Yii::$app->redis;
  102. $msg_ids = [];
  103. foreach ($msg_list as $v) {
  104. $redis->rpush($this->getMessageRedisKey(), Json::encode($v));
  105. $msg_ids[] = $v['id'];
  106. }
  107. $this->api->removeMessagePoolByParam($msg_ids);
  108. } catch (Exception $e) {
  109. Yii::error(
  110. "删除消息失败ID为" . Json::encode($msg_ids) . ", 失败内容为:" . $e->getMessage()
  111. );
  112. }
  113. }
  114. /**
  115. * 获取消息插入数据库
  116. *
  117. * @return void
  118. */
  119. public function getMessage()
  120. {
  121. Yii::$app->services->rabbitMq->push(
  122. RabbitMqEnum::EXCHANGE_CHAIN, RabbitMqEnum::CHAIN_QUEUE, ['mall_id' => $this->mall_id], self::class, 'handleSyncMsg'
  123. );
  124. }
  125. /**
  126. * 异步执行
  127. *
  128. * @return void
  129. */
  130. public static function handleSyncMsg($data)
  131. {
  132. $service = new static(['mall_id' => $data['mall_id']]);
  133. $service->apiLoadMsgToDB();
  134. }
  135. /**
  136. * api获取消息并存储
  137. *
  138. * @return void
  139. */
  140. public function apiLoadMsgToDB()
  141. {
  142. // 未安装
  143. if (empty($this->api)) return false;
  144. // 全部消息
  145. try {
  146. $msg_list = $this->api->messagePool();
  147. } catch (Exception $e) {
  148. Yii::error(
  149. "获取消息失败,mall" . $this->mall_id . ", 失败内容为:" . $e->getMessage()
  150. );
  151. }
  152. if (empty($msg_list)) return false;
  153. // 消息ID
  154. $msg_ids = [];
  155. try {
  156. // 写入数据库
  157. Message::addMsg($this->mall_id, $msg_list);
  158. // 加入队列
  159. foreach ($msg_list as $item) {
  160. $item['mall_id'] = $this->mall_id;
  161. // 延迟执行
  162. Yii::$app->services->rabbitMq->delayTask(
  163. static::DELAY, $item, self::class, 'handleSync', 'chain'
  164. );
  165. $msg_ids[] = $item['id'];
  166. }
  167. // 删除消息
  168. $this->api->removeMessagePoolByParam($msg_ids);
  169. } catch (Exception $e) {
  170. Yii::error('星链供应链消息存储失败');
  171. Yii::error($msg_list);
  172. }
  173. }
  174. /**
  175. * 获取消息池消息(redis消息队列)
  176. *
  177. * @return void
  178. */
  179. public function messagePoolHandler()
  180. {
  181. // 获取消息
  182. $redis = Yii::$app->redis;
  183. $msg_list = $redis->eval($this->_lua_get_msg(), 2, $this->getMessageRedisKey(), self::LIMIT);
  184. if (empty($msg_list)) return false;
  185. $msg_list = array_map(function ($item) {
  186. return Json::decode($item);
  187. }, $msg_list);
  188. // 储存消息
  189. Message::addMsg($this->mall_id, $msg_list);
  190. $msg_ids = array_column($msg_list, 'id');
  191. $msg_list = $this->groupMsg($msg_list);
  192. $this->messageListHandle($msg_list);
  193. }
  194. /**
  195. * 消息处理
  196. *
  197. * @param array $data
  198. * @return void
  199. */
  200. public static function handleSync(array $data)
  201. {
  202. $service = new static(['mall_id' => $data['mall_id']]);
  203. $msg_list = $service->groupMsg([$data]);
  204. $service->messageListHandle($msg_list);
  205. }
  206. /**
  207. * 处理消息
  208. *
  209. * @param array $msg_list
  210. * @return void
  211. */
  212. public function messageListHandle($msg_list)
  213. {
  214. // $db = Yii::$app->db->beginTransaction();
  215. try {
  216. foreach ($msg_list as $enum => $v) {
  217. // echo "消息类型{$dasprid-enum}, 数量" . count($v);
  218. // 按类型进行消息处理
  219. if ($func = $this->ex_func[$enum]) {
  220. // 处理消息
  221. $this->$func($v);
  222. }
  223. }
  224. // 完成处理
  225. $this->finish();
  226. // $db->commit();
  227. } catch (Exception $e) {
  228. Yii::error('星链供应链消息处理失败: ' . json_encode([
  229. 'file' => $e->getFile(),
  230. 'msg' => $e->getMessage(),
  231. 'line' => $e->getLine(),
  232. ]));
  233. Yii::error($msg_list);
  234. // $db->rollBack();
  235. }
  236. }
  237. /**
  238. * 完成消息处理
  239. *
  240. * @return void
  241. */
  242. public function finish()
  243. {
  244. // 处理成功
  245. Message::setHandled($this->successIds);
  246. // 处理失败
  247. foreach ($this->failMsgs as $fail_msg) {
  248. Message::setHandledFail($fail_msg['ids'], $fail_msg['msg']);
  249. }
  250. $this->clearMsgs();
  251. }
  252. /**
  253. * 消息按内容分组
  254. *
  255. * @param array $msgs 消息数组
  256. * @return array
  257. */
  258. public function groupMsg($msgs) :array
  259. {
  260. $group_msg = [];
  261. foreach ($msgs as $v) {
  262. $group_msg[$v['enumNum']][] = $v;
  263. }
  264. return $group_msg;
  265. }
  266. /**
  267. * 选品移除
  268. *
  269. * @param array $msg_list
  270. * @return void
  271. */
  272. public function selectionRemove($msg_list)
  273. {
  274. $service = new SpuService(['mall_id' => $this->mall_id]);
  275. foreach ($msg_list as $v) {
  276. // 更新的SPU_ID
  277. $content = Json::decode($v['content']);
  278. $spu_ids[] = strval($content['spuId']);
  279. }
  280. $service->removeImportSpuIds($spu_ids);
  281. }
  282. /**
  283. * 添加选品
  284. *
  285. * @param array $msg_list
  286. * @return void
  287. */
  288. public function selectionAdd($msg_list)
  289. {
  290. // 添加选品不做操作
  291. }
  292. /**
  293. * 商品下架
  294. *
  295. * @param array $msg_list
  296. * @return void
  297. */
  298. public function goodsPulled($msg_list)
  299. {
  300. $ids = [];
  301. try {
  302. $spu_ids = [];
  303. foreach ($msg_list as $v) {
  304. // 更新的SPU_ID
  305. $content = Json::decode($v['content']);
  306. $spu_ids[] = strval($content['spuId']);
  307. $ids[] = $v['id'];
  308. }
  309. // 更新商城商品数据
  310. $service = new GoodsService(['mall_id' => $this->mall_id]);
  311. $service->goodsPulled($spu_ids);
  312. } catch (Exception $e) {
  313. return $this->addFailMsgs($ids, $e->getMessage());
  314. }
  315. return $this->addSuccessIds($ids);
  316. }
  317. /**
  318. * 订单发货
  319. *
  320. * @param array $msg_list
  321. * @return array
  322. */
  323. public function orderDelivery($msg_list)
  324. {
  325. // 使用api查询方式
  326. if (!$this->msgDeliveryType) {
  327. $ids = [];
  328. $order_sn = [];
  329. foreach ($msg_list as $v) {
  330. // 更新的SPU_ID
  331. $content = Json::decode($v['content']);
  332. $order_sn[] = strval($content['orderSn']);
  333. $ids[] = $v['id'];
  334. }
  335. try {
  336. $service = new OrderService(['mall_id' => $this->mall_id]);
  337. $service->orderDeliveryByOrderSn($order_sn);
  338. } catch (Exception $e) {
  339. return $this->addFailMsgs($ids, $e->getMessage());
  340. }
  341. return $this->addSuccessIds($ids);
  342. }
  343. // 直接使用推送数据
  344. if ($this->msgDeliveryType) {
  345. $service = new OrderService(['mall_id' => $this->mall_id]);
  346. foreach ($msg_list as $v) {
  347. // 推送订单数据
  348. $content = Json::decode($v['content']);
  349. $order_sn = strval($content['orderSn']);
  350. $delivery_sn = strval($content['deliverySn']);
  351. $sku_ids = $content['skuIdList'];
  352. $delivery_corp_sn = $content['deliveryCorpSn'];
  353. try {
  354. // 发货
  355. $service->orderDeliveryByPushMsg(
  356. $order_sn,
  357. $delivery_sn,
  358. $sku_ids,
  359. $delivery_corp_sn
  360. );
  361. } catch (Exception $e) {
  362. return $this->addFailMsgs([$v['id']], $e->getMessage());
  363. }
  364. return $this->addSuccessIds([$v['id']]);
  365. }
  366. }
  367. }
  368. /**
  369. * 售后订单状态变更
  370. *
  371. * @param array $msg_list
  372. * status 售后状态(
  373. * 1,待卖家审核
  374. * 2,卖家拒绝退款,
  375. * 3,退款成功,
  376. * 4,卖家拒绝退货退款,
  377. * 5,待买家退货,
  378. * 6,买家已退货,待卖家收货,
  379. * 7,买家已退货,卖家拒绝收货,
  380. * 8,卖家已收货,待确认退款,
  381. * 9,退货退款成功,
  382. * 10,卖家拒绝换货,
  383. * 11,买家已退货,待卖家换货,
  384. * 12,卖家已换货,待买家收货,
  385. * 13,换货成功,
  386. * 14,已关闭,
  387. * 15,卖家同意退款,
  388. * 16,卖家同意仅退款,
  389. * 17,卖家拒绝仅退款,
  390. * 18,仅退款成功
  391. * )
  392. * @return void
  393. */
  394. public function refundOrderChange($msg_list)
  395. {
  396. /**
  397. * @var RefundService
  398. */
  399. $service = new RefundService(['mall_id' => $this->mall_id]);
  400. foreach ($msg_list as $v) {
  401. // 消息内容
  402. $content = Json::decode($v['content']);
  403. // 修改售后信息
  404. if ($service->statusChange($content['returnSn'], $content['status'])) {
  405. $this->addSuccessIds([$v['id']]);
  406. } else {
  407. $this->addFailMsgs([$v['id']], $service->getErrorMsg());
  408. }
  409. }
  410. }
  411. /**
  412. * 订单更变事件
  413. *
  414. * @param array $msg_list
  415. * @return void
  416. */
  417. public function orderChange($msg_list)
  418. {
  419. $service = new OrderService(['mall_id' => $this->mall_id]);
  420. foreach ($msg_list as $v) {
  421. // 消息内容
  422. $content = Json::decode($v['content']);
  423. // 修改订单信息
  424. if ($service->orderStatusChange($content['orderSn'], $content['status'])) {
  425. $this->addSuccessIds([$v['id']]);
  426. } else {
  427. $this->addFailMsgs([$v['id']], $service->getErrorMsg());
  428. }
  429. }
  430. }
  431. /**
  432. * 商品信息变更处理
  433. *
  434. * @param array $msg_list
  435. * @return void
  436. */
  437. public function goodsChange($msg_list)
  438. {
  439. $ids = [];
  440. try {
  441. $spu_ids = [];
  442. foreach ($msg_list as $v) {
  443. // 更新的SPU_ID
  444. $content = Json::decode($v['content']);
  445. $spu_ids[] = strval($content['spuId']);
  446. $ids[] = $v['id'];
  447. }
  448. // 接口获取商品数据
  449. $this->api->setNotExcption();
  450. $spu_list = array_map(function ($item) {
  451. return $item[0];
  452. }, $this->api->getSpuBySpuIdsPool($spu_ids));
  453. // 更新供应链商品数据
  454. Goods::saveList($spu_list);
  455. // 更新商城商品数据
  456. $service = new GoodsService(['mall_id' => $this->mall_id]);
  457. foreach ($spu_list as $spu) {
  458. $service->editGoods($spu['spuId']);
  459. }
  460. } catch (Exception $e) {
  461. return $this->addFailMsgs($ids, $e->getMessage());
  462. }
  463. return $this->addSuccessIds($ids);
  464. }
  465. /**
  466. * 不做处理
  467. *
  468. * @param array $msg_list
  469. * @return array
  470. */
  471. public function notHandler($msg_list)
  472. {
  473. $msg_ids = array_column($msg_list, 'id');
  474. return $msg_ids;
  475. }
  476. /**
  477. * lua脚本
  478. * _key redis key
  479. * _limit 循环取出数量
  480. *
  481. * @return string
  482. */
  483. public function _lua_get_msg()
  484. {
  485. return <<<SCRIPT
  486. -- key
  487. local _key = KEYS[1]
  488. local _limit = tonumber(KEYS[2])
  489. local _count = redis.call('LLEN', _key)
  490. if _limit > _count then
  491. _limit = _count
  492. end
  493. if _count > 0 then
  494. local msg = {}
  495. -- 循环_limit
  496. for i = 1, _limit
  497. do
  498. msg[i] = redis.call('LPOP', _key)
  499. end
  500. return msg
  501. else
  502. return 0
  503. end
  504. SCRIPT;
  505. }
  506. /**
  507. * 添加成功msg_id
  508. *
  509. * @param array $ids
  510. * @return void
  511. */
  512. public function addSuccessIds(array $ids)
  513. {
  514. $this->successIds = array_merge($this->successIds, $ids);
  515. }
  516. /**
  517. * 添加处理失败消息
  518. *
  519. * @param array $ids
  520. * @param string $msg
  521. * @return void
  522. */
  523. public function addFailMsgs(array $ids, string $msg)
  524. {
  525. $this->failMsgs[] = [
  526. 'ids' => $ids,
  527. 'msg' => $msg
  528. ];
  529. }
  530. /**
  531. * 清空消息
  532. *
  533. * @return void
  534. */
  535. public function clearMsgs()
  536. {
  537. $this->successIds = [];
  538. $this->failMsgs = [];
  539. }
  540. /**
  541. * 测试数据
  542. *
  543. * @return array
  544. */
  545. public function testMsgList()
  546. {
  547. $msg_list[] = [
  548. 'id' => '66346192',
  549. 'tenantId' => '1535098852173492226',
  550. 'enumNum' => 8,
  551. 'content' => '{"spuId":1426102374728585218}',
  552. 'createTime' => '2022-07-22 14:28:59',
  553. ];
  554. // $msg_list = array_map(function ($item) {
  555. // return Json::decode($item);
  556. // }, $msg_list);
  557. $msg_list = $this->groupMsg($msg_list);
  558. $this->messageListHandle($msg_list);
  559. /* $msg_list[] = [
  560. 'id' => '63503428',
  561. 'tenantId' => '1535098852173492226',
  562. 'enumNum' => 3,
  563. 'content' => '{"origStatus":1,"origStatusName":"待卖家审核","returnSn":"SH8763722084","status":15,"statusName":"卖家同意退款"}',
  564. 'createTime' => '2022-06-27 11:49:06',
  565. ]; */
  566. /* $msg_list[] = [
  567. 'id' => '63609622',
  568. 'tenantId' => '1535098852173492226',
  569. 'enumNum' => 1,
  570. 'content' => '{"orderSn":"PO182297850412888","origStatus":20,"origStatusName":"待发货","status":30,"statusName":"待收货"}',
  571. 'createTime' => '2022-06-28 10:05:03',
  572. ];
  573. $msg_list[] = [
  574. 'id' => '63609625',
  575. 'tenantId' => '1535098852173492226',
  576. 'enumNum' => 2,
  577. 'content' => '{"orderSn":"PO182297850412888","skuIdList":["1526494051752783874"]}',
  578. 'createTime' => '2022-06-28 10:05:03',
  579. ]; */
  580. /* $msg_list[] = [
  581. 'id' => '63625249',
  582. 'tenantId' => '1535098852173492226',
  583. 'enumNum' => 3,
  584. 'content' => '{"origStatus":1,"origStatusName":"待卖家审核","returnSn":"SH4873942686","status":5,"statusName":"待买家退货"}',
  585. 'createTime' => '2022-06-28 11:27:30',
  586. ]; */
  587. }
  588. }