services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK, RabbitMqEnum::TASK_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 重试队列消费 * @throws Exception */ public function actionRetryConsume() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_DELAY_RETRY, RabbitMqEnum::DELAY_RETRY_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 订单任务队列消费 * @throws \Exception */ public function actionOrderConsume() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_ORDER, RabbitMqEnum::ORDER_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 延时队列消费 * @throws Exception */ public function actionDelayConsume() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_DELAY, RabbitMqEnum::DELAY_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 消息处理 * @param AMQPMessage $message * @throws * @return bool */ private function handleMessage($message) { $routing_key = $message->getRoutingKey(); echo '队列(' . $routing_key . ')的消息,其内容为:' . $message->getBody() . PHP_EOL . PHP_EOL; $message_body = json_decode($message->getBody(), true); $class = $message_body['handler_class']; $method = $message_body['method']; $data = $message_body['data']; if (!method_exists($class, $method)) { //方法不存在,忽略 Yii::error('actionTaskConsume: class' . $class . ' ,method:' . $method . ' ,not exist'); return false; } $reflect_method = new ReflectionMethod($class, $method); Yii::$app->db->close(); Yii::$app->redis->close(); if ($reflect_method->isStatic()) { $res = $class::$method($data); if ($res === false) Yii::error($class::$static_error ?? '未定义错误'); //if($res === false) throw new Exception($class::$static_error ?? '未定义错误'); } else { $obj = new $class(); $res = $obj->$method($data); if ($res === false) Yii::error($obj->error ?? '未定义错误'); //if($res === false) throw new Exception($obj->error ?? '未定义错误'); } return $res; } /** * 修改关系链队列消费 * @throws \Exception */ public function actionChangeRelationshipConsume() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_CHANGE_RELATIONSHIP, RabbitMqEnum::CHANGE_RELATIONSHIP_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 编辑商品队列消费 * @throws \Exception */ public function actionEditGoods() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_GOODS, RabbitMqEnum::GOODS_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 编辑用户资金队列消费 * @throws \Exception */ public function actionUserWalletConsume() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_USER_WALLET, RabbitMqEnum::USER_WALLET_QUEUE, function ($message) { $message_id = $message->get('message_id'); $data['message_id'] = $message->get('message_id'); $routing_key = $message->getRoutingKey(); echo $message->getConsumerTag() . '接收到发给队列(' . $routing_key . ')的消息,其内容为:' . $message->getBody() . PHP_EOL; $message_body = json_decode($message->getBody(), true); $class = $message_body['handler_class']; $method = $message_body['method']; $data = $message_body['data']; $res = $class::$method($data); //$res = UserWalletService::handle($message_id,$data); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'actionUserWalletConsume exception:' . $e->getMessage(); } } /** * 拼团任务队列消费 * @throws \Exception */ public function actionGroupBuy() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_GROUP_BUY, RabbitMqEnum::GROUP_BUY_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 幸运拼团任务队列消费 * @throws \Exception */ public function actionLuckyGroup() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_LUCKY_GROUP, RabbitMqEnum::LUCKY_GROUP_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 幸运拼团任务队列消费 * @throws \Exception */ public function actionCensus() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_CENSUS, RabbitMqEnum::CENSUS_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } public function actionImportUser() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_TASK_IMPORT_USER, RabbitMqEnum::IMPORT_USER, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } /** * 长耗时队列消费 * @throws \Exception */ public function actionLongConsume() { try { Yii::$app->services->rabbitMq->listen(RabbitMqEnum::EXCHANGE_LONG_DURATION, RabbitMqEnum::LONG_DURATION_QUEUE, function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage(); } } }