RabbitMqEnum::EXCHANGE_TASK, 'queue_name' => RabbitMqEnum::TASK_QUEUE, 'worker_num' => 1, 'exchange_type' => AMQPExchangeType::DIRECT, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_ORDER, 'queue_name' => RabbitMqEnum::ORDER_QUEUE, 'worker_num' => 1, 'exchange_type' => AMQPExchangeType::DIRECT, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_TASK_USER_WALLET, 'queue_name' => RabbitMqEnum::USER_WALLET_QUEUE, 'worker_num' => 1, 'exchange_type' => AMQPExchangeType::DIRECT ], [ 'exchange' => RabbitMqEnum::EXCHANGE_TASK_CHANGE_RELATIONSHIP, 'queue_name' => RabbitMqEnum::CHANGE_RELATIONSHIP_QUEUE, 'worker_num' => 1, 'exchange_type' => AMQPExchangeType::DIRECT ], [ 'exchange' => RabbitMqEnum::EXCHANGE_DELAY_RETRY, 'queue_name' => RabbitMqEnum::DELAY_RETRY_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_TASK_GROUP_BUY, 'queue_name' => RabbitMqEnum::GROUP_BUY_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_TASK_CENSUS, 'queue_name' => RabbitMqEnum::CENSUS_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_TASK_LUCKY_GROUP, 'queue_name' => RabbitMqEnum::LUCKY_GROUP_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_DELAY, 'queue_name' => RabbitMqEnum::DELAY_HANDLE_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_STATISTICS, 'queue_name' => RabbitMqEnum::STATISTICS_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_TASK_IMPORT_USER, 'queue_name' => RabbitMqEnum::IMPORT_USER, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_CHAIN, 'queue_name' => RabbitMqEnum::CHAIN_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_UPLOAD, 'queue_name' => RabbitMqEnum::UPLOAD_QUEUE, 'worker_num' => 2, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_CAKE, 'queue_name' => RabbitMqEnum::CAKE_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_CAKE_ORDER, 'queue_name' => RabbitMqEnum::CAKE_ORDER_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_LONG_DURATION, 'queue_name' => RabbitMqEnum::LONG_DURATION_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_QSC_ORDER, 'queue_name' => RabbitMqEnum::QSC_ORDER_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_QSC, 'queue_name' => RabbitMqEnum::QSC_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_CARD_COUPONS, 'queue_name' => RabbitMqEnum::CARD_COUPONS_QUEUE, 'worker_num' => 1, ], [ 'exchange' => RabbitMqEnum::EXCHANGE_LONG_ORDER, 'queue_name' => RabbitMqEnum::LONG_ORDER_QUEUE, 'worker_num' => 2, ], ]; // 进程与消费者的对应关系 public $worker_queue_relate = []; public function init(){ $worker_num_arr = array_column($this->consume_setting,'worker_num'); $this->worker_num = array_sum($worker_num_arr); $worker_id = 0; foreach($this->consume_setting as $key=>$queue_info){ $wait_assign_worker_num = $queue_info['worker_num']; while($wait_assign_worker_num > 0){ $this->worker_queue_relate[$worker_id] = $key; $worker_id ++; $wait_assign_worker_num --; } } } /** * 启动消费者 * @Author bing * @DateTime 2021-10-20 16:23:03 * @copyright: Copyright (c) 2020 广东七件事集团 * @return void */ public function actionRun(){ // 计数器,当进程反复退出重启时,可能代码有致命错误,需要处理 // $atomic = new Atomic(); $pool = new Pool($this->worker_num); //绑定一个事件 $pool->on("WorkerStart", function ($pool, $worker_id) { echo 'worker_id:'.$worker_id.PHP_EOL; $queue_key = $this->worker_queue_relate[$worker_id]; $queue_info = $this->consume_setting[$queue_key]; //设置进程名称 swoole_set_process_name('swoole : rabbitmq_process '.$queue_info['queue_name'].'_'.$worker_id); $process = $pool->getProcess($worker_id); try { Yii::$app->services->rabbitMq->listen($queue_info['exchange'], $queue_info['queue_name'], function ($message) { $res = $this->handleMessage($message); Yii::getLogger()->flush(true); return $res; }); } catch (Exception $e) { Yii::getLogger()->flush(true); echo 'exception:' . $e->getMessage().$e->getFile().':'.$e->getLine().PHP_EOL; } }); //子进程关闭 $pool->on("WorkerStop", function ($pool, $workerId) { echo "Worker#{$workerId} is stopped\n"; }); $pool->start(); } /** * 消息处理 * @param AMQPMessage $message * @throws * @return bool */ private function handleMessage($message) { $date = date('Y-m-d H:i:s'); $routing_key = $message->getRoutingKey(); $flush = '['.$date.']'.'队列(' . $routing_key . ')接收到消息,其内容为:' . $message->getBody(); $message_body = json_decode($message->getBody(), true); $class = $message_body['handler_class']; $method = $message_body['method']; $data = $message_body['data']; $data['message_id'] = $message->get('message_id'); 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(); $time1 = time(); 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 ?? '未定义错误'); } $time2 = time(); $total_time = $time2 - $time1; if ($total_time > 1) $flush .= ',队列耗时:' .$total_time . '秒'; echo $flush . PHP_EOL . PHP_EOL; return $res; } /** * 重启消费者 * @Author bing * @DateTime 2021-10-20 16:23:03 * @copyright: Copyright (c) 2020 广东七件事集团 * @return void */ public function actionRestartConsume(){ } }