SwooleConsumeController.php 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240
  1. <?php
  2. /**
  3. * 进程池多消费者,多队列消费rabbitMQ队列
  4. */
  5. namespace console\controllers;
  6. use common\enums\RabbitMqEnum;
  7. use console\controllers\swoole_server\BaseSwooleServer;
  8. use Exception;
  9. use PhpAmqpLib\Exchange\AMQPExchangeType;
  10. use ReflectionMethod;
  11. use Swoole\Process\Pool;
  12. use Yii;
  13. /**
  14. * rabbitMQ队列消费者进程池
  15. * @Author bing
  16. * @DateTime 2021-10-21 10:30:19
  17. * @copyright: Copyright (c) 2020 广东七件事集团
  18. */
  19. class SwooleConsumeController extends BaseSwooleServer
  20. {
  21. public $worker_num;
  22. public $consume_setting = [
  23. [
  24. 'exchange' => RabbitMqEnum::EXCHANGE_TASK,
  25. 'queue_name' => RabbitMqEnum::TASK_QUEUE,
  26. 'worker_num' => 1,
  27. 'exchange_type' => AMQPExchangeType::DIRECT,
  28. ],
  29. [
  30. 'exchange' => RabbitMqEnum::EXCHANGE_ORDER,
  31. 'queue_name' => RabbitMqEnum::ORDER_QUEUE,
  32. 'worker_num' => 1,
  33. 'exchange_type' => AMQPExchangeType::DIRECT,
  34. ],
  35. [
  36. 'exchange' => RabbitMqEnum::EXCHANGE_TASK_USER_WALLET,
  37. 'queue_name' => RabbitMqEnum::USER_WALLET_QUEUE,
  38. 'worker_num' => 1,
  39. 'exchange_type' => AMQPExchangeType::DIRECT
  40. ],
  41. [
  42. 'exchange' => RabbitMqEnum::EXCHANGE_TASK_CHANGE_RELATIONSHIP,
  43. 'queue_name' => RabbitMqEnum::CHANGE_RELATIONSHIP_QUEUE,
  44. 'worker_num' => 1,
  45. 'exchange_type' => AMQPExchangeType::DIRECT
  46. ],
  47. [
  48. 'exchange' => RabbitMqEnum::EXCHANGE_DELAY_RETRY,
  49. 'queue_name' => RabbitMqEnum::DELAY_RETRY_QUEUE,
  50. 'worker_num' => 1,
  51. ],
  52. [
  53. 'exchange' => RabbitMqEnum::EXCHANGE_TASK_GROUP_BUY,
  54. 'queue_name' => RabbitMqEnum::GROUP_BUY_QUEUE,
  55. 'worker_num' => 1,
  56. ],
  57. [
  58. 'exchange' => RabbitMqEnum::EXCHANGE_TASK_CENSUS,
  59. 'queue_name' => RabbitMqEnum::CENSUS_QUEUE,
  60. 'worker_num' => 1,
  61. ],
  62. [
  63. 'exchange' => RabbitMqEnum::EXCHANGE_TASK_LUCKY_GROUP,
  64. 'queue_name' => RabbitMqEnum::LUCKY_GROUP_QUEUE,
  65. 'worker_num' => 1,
  66. ],
  67. [
  68. 'exchange' => RabbitMqEnum::EXCHANGE_DELAY,
  69. 'queue_name' => RabbitMqEnum::DELAY_HANDLE_QUEUE,
  70. 'worker_num' => 1,
  71. ],
  72. [
  73. 'exchange' => RabbitMqEnum::EXCHANGE_STATISTICS,
  74. 'queue_name' => RabbitMqEnum::STATISTICS_QUEUE,
  75. 'worker_num' => 1,
  76. ],
  77. [
  78. 'exchange' => RabbitMqEnum::EXCHANGE_TASK_IMPORT_USER,
  79. 'queue_name' => RabbitMqEnum::IMPORT_USER,
  80. 'worker_num' => 1,
  81. ],
  82. [
  83. 'exchange' => RabbitMqEnum::EXCHANGE_CHAIN,
  84. 'queue_name' => RabbitMqEnum::CHAIN_QUEUE,
  85. 'worker_num' => 1,
  86. ],
  87. [
  88. 'exchange' => RabbitMqEnum::EXCHANGE_UPLOAD,
  89. 'queue_name' => RabbitMqEnum::UPLOAD_QUEUE,
  90. 'worker_num' => 2,
  91. ],
  92. [
  93. 'exchange' => RabbitMqEnum::EXCHANGE_CAKE,
  94. 'queue_name' => RabbitMqEnum::CAKE_QUEUE,
  95. 'worker_num' => 1,
  96. ],
  97. [
  98. 'exchange' => RabbitMqEnum::EXCHANGE_CAKE_ORDER,
  99. 'queue_name' => RabbitMqEnum::CAKE_ORDER_QUEUE,
  100. 'worker_num' => 1,
  101. ],
  102. [
  103. 'exchange' => RabbitMqEnum::EXCHANGE_LONG_DURATION,
  104. 'queue_name' => RabbitMqEnum::LONG_DURATION_QUEUE,
  105. 'worker_num' => 1,
  106. ],
  107. [
  108. 'exchange' => RabbitMqEnum::EXCHANGE_QSC_ORDER,
  109. 'queue_name' => RabbitMqEnum::QSC_ORDER_QUEUE,
  110. 'worker_num' => 1,
  111. ],
  112. [
  113. 'exchange' => RabbitMqEnum::EXCHANGE_QSC,
  114. 'queue_name' => RabbitMqEnum::QSC_QUEUE,
  115. 'worker_num' => 1,
  116. ],
  117. [
  118. 'exchange' => RabbitMqEnum::EXCHANGE_CARD_COUPONS,
  119. 'queue_name' => RabbitMqEnum::CARD_COUPONS_QUEUE,
  120. 'worker_num' => 1,
  121. ],
  122. [
  123. 'exchange' => RabbitMqEnum::EXCHANGE_LONG_ORDER,
  124. 'queue_name' => RabbitMqEnum::LONG_ORDER_QUEUE,
  125. 'worker_num' => 2,
  126. ],
  127. ];
  128. // 进程与消费者的对应关系
  129. public $worker_queue_relate = [];
  130. public function init(){
  131. $worker_num_arr = array_column($this->consume_setting,'worker_num');
  132. $this->worker_num = array_sum($worker_num_arr);
  133. $worker_id = 0;
  134. foreach($this->consume_setting as $key=>$queue_info){
  135. $wait_assign_worker_num = $queue_info['worker_num'];
  136. while($wait_assign_worker_num > 0){
  137. $this->worker_queue_relate[$worker_id] = $key;
  138. $worker_id ++;
  139. $wait_assign_worker_num --;
  140. }
  141. }
  142. }
  143. /**
  144. * 启动消费者
  145. * @Author bing
  146. * @DateTime 2021-10-20 16:23:03
  147. * @copyright: Copyright (c) 2020 广东七件事集团
  148. * @return void
  149. */
  150. public function actionRun(){
  151. // 计数器,当进程反复退出重启时,可能代码有致命错误,需要处理
  152. // $atomic = new Atomic();
  153. $pool = new Pool($this->worker_num);
  154. //绑定一个事件
  155. $pool->on("WorkerStart", function ($pool, $worker_id) {
  156. echo 'worker_id:'.$worker_id.PHP_EOL;
  157. $queue_key = $this->worker_queue_relate[$worker_id];
  158. $queue_info = $this->consume_setting[$queue_key];
  159. //设置进程名称
  160. swoole_set_process_name('swoole : rabbitmq_process '.$queue_info['queue_name'].'_'.$worker_id);
  161. $process = $pool->getProcess($worker_id);
  162. try {
  163. Yii::$app->services->rabbitMq->listen($queue_info['exchange'], $queue_info['queue_name'], function ($message) {
  164. $res = $this->handleMessage($message);
  165. Yii::getLogger()->flush(true);
  166. return $res;
  167. });
  168. } catch (Exception $e) {
  169. Yii::getLogger()->flush(true);
  170. echo 'exception:' . $e->getMessage().$e->getFile().':'.$e->getLine().PHP_EOL;
  171. }
  172. });
  173. //子进程关闭
  174. $pool->on("WorkerStop", function ($pool, $workerId) {
  175. echo "Worker#{$workerId} is stopped\n";
  176. });
  177. $pool->start();
  178. }
  179. /**
  180. * 消息处理
  181. * @param AMQPMessage $message
  182. * @throws
  183. * @return bool
  184. */
  185. private function handleMessage($message)
  186. {
  187. $date = date('Y-m-d H:i:s');
  188. $routing_key = $message->getRoutingKey();
  189. $flush = '['.$date.']'.'队列(' . $routing_key . ')接收到消息,其内容为:' . $message->getBody();
  190. $message_body = json_decode($message->getBody(), true);
  191. $class = $message_body['handler_class'];
  192. $method = $message_body['method'];
  193. $data = $message_body['data'];
  194. $data['message_id'] = $message->get('message_id');
  195. if (!method_exists($class, $method)) {
  196. //方法不存在,忽略
  197. Yii::error('actionTaskConsume: class' . $class . ' ,method:' . $method . ' ,not exist');
  198. return false;
  199. }
  200. $reflect_method = new ReflectionMethod($class, $method);
  201. Yii::$app->db->close();
  202. Yii::$app->redis->close();
  203. $time1 = time();
  204. if ($reflect_method->isStatic()) {
  205. $res = $class::$method($data);
  206. if ($res === false) Yii::error($class::$static_error ?? '未定义错误');
  207. //if($res === false) throw new Exception($class::$static_error ?? '未定义错误');
  208. } else {
  209. $obj = new $class();
  210. $res = $obj->$method($data);
  211. if ($res === false) Yii::error($obj->error ?? '未定义错误');
  212. //if($res === false) throw new Exception($obj->error ?? '未定义错误');
  213. }
  214. $time2 = time();
  215. $total_time = $time2 - $time1;
  216. if ($total_time > 1) $flush .= ',队列耗时:' .$total_time . '秒';
  217. echo $flush . PHP_EOL . PHP_EOL;
  218. return $res;
  219. }
  220. /**
  221. * 重启消费者
  222. * @Author bing
  223. * @DateTime 2021-10-20 16:23:03
  224. * @copyright: Copyright (c) 2020 广东七件事集团
  225. * @return void
  226. */
  227. public function actionRestartConsume(){
  228. }
  229. }