RabbitMqService.php 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493
  1. <?php
  2. namespace services\common;
  3. use common\components\Service;
  4. use PhpAmqpLib\Connection\AMQPStreamConnection;
  5. use PhpAmqpLib\Exchange\AMQPExchangeType;
  6. use PhpAmqpLib\Message\AMQPMessage;
  7. use PhpAmqpLib\Wire\AMQPTable;
  8. use Exception;
  9. use Closure;
  10. use common\enums\RabbitMqEnum;
  11. use Yii;
  12. use yii\redis\Connection;
  13. /**
  14. * RabbitMq服务类
  15. * Date: 2021/8/20
  16. * Time: 14:56
  17. */
  18. class RabbitMqService extends Service
  19. {
  20. public $host;
  21. public $port = 5672;
  22. public $userName;
  23. public $passWord;
  24. public $setting = [
  25. 'vhost' => '/',
  26. 'insist' => false,
  27. 'login_method' => 'AMQPLAIN',
  28. 'login_response' => null,
  29. 'locale' => 'en_US',
  30. 'connection_timeout' => 3.0,
  31. 'read_write_timeout' => 3.0,
  32. 'context' => null,
  33. 'keepalive' => false,
  34. 'heartbeat' => 0
  35. ];
  36. // Send a message with the string "quit" to cancel the consumer.
  37. public $cancelTag = 'quit';
  38. public $connection = null;
  39. /**
  40. * 获取一个rabbitMQ连接
  41. * @Author bing
  42. * @DateTime 2021-02-24 10:35:37
  43. * @copyright: Copyright (c) 2020 广东七件事集团
  44. * @return AMQPStreamConnection
  45. * @throws Exception
  46. */
  47. public function connect()
  48. {
  49. $rabbitMqConfig = \Yii::$app->getComponents()['rabbitMq'];
  50. unset($rabbitMqConfig['class']);
  51. if ($rabbitMqConfig) {
  52. $this->host = $rabbitMqConfig['host'];
  53. $this->port = $rabbitMqConfig['port'];
  54. $this->userName = $rabbitMqConfig['userName'];
  55. $this->passWord = $rabbitMqConfig['passWord'];
  56. $this->setting = $rabbitMqConfig['setting'];
  57. }
  58. $config = [
  59. 'host' => $this->host,
  60. 'port' => $this->port,
  61. 'user' => $this->userName,
  62. 'password' => $this->passWord,
  63. 'vhost' => $this->setting['vhost'],
  64. 'insist' => $this->setting['insist'],
  65. 'login_method' => $this->setting['login_method'],
  66. 'login_response' => $this->setting['login_response'],
  67. 'locale' => $this->setting['locale'],
  68. 'connection_timeout' => $this->setting['connection_timeout'],
  69. 'read_write_timeout' => $this->setting['read_write_timeout'],
  70. 'context' => $this->setting['context'],
  71. 'keepalive' => $this->setting['keepalive'],
  72. 'heartbeat'=>$this->setting['heartbeat']
  73. ];
  74. $connection = new AMQPStreamConnection(...array_values($config));
  75. if (!$connection->isConnected()) throw new Exception('Connect failed!');
  76. $this->connection = $connection;
  77. return $connection;
  78. }
  79. /**
  80. * 生产一个消息并发送到指定(direct)交换机
  81. * @Author bing
  82. * @DateTime 2021-02-24 15:28:47
  83. * @copyright: Copyright (c) 2020 广东七件事集团
  84. * @param string $exchange
  85. * @param string $queue_name 队列名称
  86. * @param array $data
  87. * @param string $handler_class
  88. * @param string $method
  89. * @param Closure $error_callback
  90. * @param Closure $success_callback
  91. * @param string $exchange_type 交换机类型:默认直连交换机
  92. * @throws Exception
  93. * @return int
  94. */
  95. public function push(string $exchange, string $queue_name, array $data, string $handler_class = '', string $method = '', Closure $error_callback = null, Closure $success_callback = null,$exchange_type = AMQPExchangeType::DIRECT)
  96. {
  97. $message['data'] = $data;
  98. $message['handler_class'] = $handler_class;
  99. $message['method'] = $method;
  100. /** @var AMQPStreamConnection $conn */
  101. $conn = $this->connect();
  102. $channel = $conn->channel();
  103. $channel->exchange_declare($exchange, $exchange_type, false, false, false);
  104. $channel->queue_declare($queue_name, false, true, false, false);
  105. //绑定队列到交换机
  106. $channel->queue_bind($queue_name, $exchange, $queue_name);
  107. $message_id = $this->createMessageId();
  108. $body = new AMQPMessage(json_encode($message,JSON_UNESCAPED_UNICODE), array('delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT, 'message_id' => $message_id, 'application_headers' => new AMQPTable(['retry' => 0])));
  109. //消息发送状态回调
  110. $channel->set_ack_handler(function (AMQPMessage $message) use ($success_callback) {
  111. !empty($success_callback) && $success_callback($message);
  112. });
  113. $channel->set_nack_handler(function (AMQPMessage $message) use ($error_callback) {
  114. !empty($error_callback) && $error_callback($message);
  115. });
  116. //开启消息发送状态监听
  117. $channel->confirm_select();
  118. $res = $this->setMessage($message_id, json_encode($message));
  119. if (!$res) {
  120. throw new Exception('设置消息失败');
  121. }
  122. $channel->basic_publish($body, $exchange, $queue_name);
  123. $channel->wait_for_pending_acks();
  124. $channel->close();
  125. $conn->close();
  126. Yii::$app->custom->rabbitMqLogs("rabbitmq-push 消息队列投递成功。投递ID:" . $message_id, json_encode($message));
  127. return $message_id;
  128. }
  129. /**
  130. * 延时队列生产
  131. * @notice:此延时队列不能保证时间的精准,当业务处理出现阻塞,则在队列里已达到过期时间的消息并不会被发送到对应的队列。
  132. * @Author bing
  133. * @DateTime 2021-02-24 16:49:31
  134. * @copyright: Copyright (c) 2020 广东七件事集团
  135. * @param int $sec 延时秒数
  136. * @param array $data 传递的数据
  137. * @param string $handler_class 处理类
  138. * @param string $method 处理方法
  139. * @param array $params 处理类参数
  140. * @throws Exception
  141. * @return void
  142. */
  143. public function delay(int $sec, array $data, string $handler_class = '', string $method = '', array $params = [])
  144. {
  145. $message['data'] = $data;
  146. $message['handler_class'] = $handler_class;
  147. $message['method'] = $method;
  148. $message['params'] = $params;
  149. $micro_sec = $sec * 1000;
  150. /** @var AMQPStreamConnection $conn */
  151. $conn = $this->connect();
  152. $channel = $conn->channel();
  153. $channel->exchange_declare(RabbitMqEnum::EXCHANGE_DELAY, 'direct', false, false, false);
  154. $channel->exchange_declare(RabbitMqEnum::EXCHANGE_DELAY_CACHE, 'direct', false, false, false);
  155. //死信交换机和路由
  156. $tale = new AMQPTable();
  157. $tale->set('x-dead-letter-exchange', RabbitMqEnum::EXCHANGE_DELAY);
  158. $tale->set('x-dead-letter-routing-key', RabbitMqEnum::DELAY_HANDLE_QUEUE);
  159. $tale->set('x-message-ttl', $micro_sec);
  160. $queue_name = 'delay_cache_queue_' . $sec . 's';
  161. //延时缓存队列声明及绑定
  162. $channel->queue_declare($queue_name, false, true, false, true, false, $tale);
  163. $channel->queue_bind($queue_name, RabbitMqEnum::EXCHANGE_DELAY_CACHE, $queue_name);
  164. //延时处理队列声明及绑定
  165. $channel->queue_declare(RabbitMqEnum::DELAY_HANDLE_QUEUE, false, true, false, false, false);
  166. $channel->queue_bind(RabbitMqEnum::DELAY_HANDLE_QUEUE, RabbitMqEnum::EXCHANGE_DELAY, RabbitMqEnum::DELAY_HANDLE_QUEUE);
  167. $message_id = $this->createMessageId();
  168. $res = $this->setMessage($message_id, json_encode($message),$sec+3600);
  169. if (!$res) {
  170. throw new Exception('设置消息失败');
  171. }
  172. $body = new AMQPMessage(json_encode($message), array(
  173. 'expiration' => $micro_sec,
  174. 'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
  175. 'message_id' => $message_id
  176. ));
  177. $channel->basic_publish($body, RabbitMqEnum::EXCHANGE_DELAY_CACHE, $queue_name);
  178. $channel->close();
  179. $conn->close();
  180. Yii::$app->custom->rabbitMqLogs("rabbitmq-delay 消息队列投递成功。投递ID:" . $message_id, json_encode($message));
  181. }
  182. /**
  183. * 延时队列生产
  184. * @notice:此延时队列不能保证时间的精准,当业务处理出现阻塞,则在队列里已达到过期时间的消息并不会被发送到对应的队列。
  185. * @Author bing
  186. * @DateTime 2021-02-24 16:49:31
  187. * @copyright: Copyright (c) 2020 广东七件事集团
  188. * @param int $sec 延时秒数
  189. * @param array $data 传递的数据
  190. * @param string $handler_class 处理类
  191. * @param string $method 处理方法
  192. * @param array $params 处理类参数
  193. * @throws Exception
  194. * @return void
  195. */
  196. public function delayTask(int $sec, array $data, string $handler_class = '', string $method = '', $consume_type_key = 'delay', array $params = [])
  197. {
  198. $message['data'] = $data;
  199. $message['handler_class'] = $handler_class;
  200. $message['method'] = $method;
  201. $message['params'] = $params;
  202. $micro_sec = $sec * 1000;
  203. $consume_type = RabbitMqEnum::getConsumeType($consume_type_key);
  204. /** @var AMQPStreamConnection $conn */
  205. $conn = $this->connect();
  206. $channel = $conn->channel();
  207. $channel->exchange_declare($consume_type['exchange'], 'direct', false, false, false);
  208. $channel->exchange_declare(RabbitMqEnum::EXCHANGE_DELAY_CACHE, 'direct', false, false, false);
  209. //死信交换机和路由
  210. $tale = new AMQPTable();
  211. $tale->set('x-dead-letter-exchange', $consume_type['exchange']);
  212. $tale->set('x-dead-letter-routing-key', $consume_type['queue_name']);
  213. $tale->set('x-message-ttl', $micro_sec);
  214. $queue_name = $consume_type_key == 'delay' ? 'delay_cache_queue_' . $sec . 's' : 'delay_'.$consume_type_key.'_' . $sec . 's';
  215. //延时缓存队列声明及绑定
  216. $channel->queue_declare($queue_name, false, true, false, true, false, $tale);
  217. $channel->queue_bind($queue_name, RabbitMqEnum::EXCHANGE_DELAY_CACHE, $queue_name);
  218. //延时处理队列声明及绑定
  219. $channel->queue_declare($consume_type['queue_name'], false, true, false, false, false);
  220. $channel->queue_bind($consume_type['queue_name'], $consume_type['exchange'], $consume_type['queue_name']);
  221. $message_id = $this->createMessageId();
  222. $res = $this->setMessage($message_id, json_encode($message),$sec+3600);
  223. if (!$res) {
  224. throw new Exception('设置消息失败');
  225. }
  226. $body = new AMQPMessage(json_encode($message), array(
  227. 'expiration' => $micro_sec,
  228. 'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
  229. 'message_id' => $message_id
  230. ));
  231. $channel->basic_publish($body, RabbitMqEnum::EXCHANGE_DELAY_CACHE, $queue_name);
  232. $channel->close();
  233. $conn->close();
  234. Yii::$app->custom->rabbitMqLogs("rabbitmq-delay-task 消息队列投递成功。投递ID:" . $message_id, json_encode($message));
  235. }
  236. /**
  237. * 消费者监听
  238. * @Author bing
  239. * @DateTime 2021-02-24 15:18:08
  240. * @copyright: Copyright (c) 2020 广东七件事集团
  241. * @param string $exchange
  242. * @param string $queue_name
  243. * @param Closure $callback
  244. * @param string $exchange_type 交换机类型:默认直连交换机
  245. * @throws Exception
  246. * @return void
  247. */
  248. public function listen(string $exchange, string $queue_name, Closure $callback = null,$exchange_type = AMQPExchangeType::DIRECT)
  249. {
  250. /** @var AMQPStreamConnection $conn */
  251. $conn = $this->connect();
  252. if (!$conn) {
  253. throw new Exception('Connect failed!');
  254. }
  255. $channel = $conn->channel();
  256. $channel->exchange_declare($exchange, $exchange_type, false, false, false);
  257. $channel->queue_declare($queue_name, false, true, false, false);
  258. //绑定队列到交换机
  259. $channel->queue_bind($queue_name, $exchange, $queue_name);
  260. $channel->basic_qos(null, 1, null);
  261. $channel->basic_consume($queue_name, '', false, false, false, false, function (AMQPMessage $message) use ($callback, $channel, $exchange, $queue_name) : void {
  262. //消息确认
  263. $message->getChannel()->basic_ack($message->getDeliveryTag());
  264. //消息判断是否消费过
  265. $message_id = $message->get('message_id');
  266. //echo 'message_id:' . $message_id.PHP_EOL;
  267. $message_content = $this->checkMessage($message_id);
  268. if (!empty($message_content) && !empty($callback)) {
  269. $res = $callback($message);
  270. if (false === $res) {
  271. //投递到 重试队列
  272. $this->retry($message);
  273. } else {
  274. //消费成功,删除消息
  275. $this->deleteMessage($message_id);
  276. }
  277. }
  278. //!empty($message_content) && !empty($callback) && $callback($message);
  279. });
  280. while (count($channel->callbacks)) {
  281. $channel->wait();
  282. }
  283. $channel->close();
  284. $conn->close();
  285. }
  286. /**
  287. * 取出一个队列里的消息
  288. * @Author bing
  289. * @DateTime 2021-02-27 17:19:25
  290. * @copyright: Copyright (c) 2020 广东七件事集团
  291. * @param string $queue_name
  292. * @return AMQPMessage|null
  293. * @throws
  294. */
  295. public function pull(string $queue_name)
  296. {
  297. /** @var AMQPStreamConnection $conn */
  298. $conn = $this->connect();
  299. $channel = $conn->channel();
  300. $channel->queue_declare($queue_name, false, true, false, false);
  301. /* @var AMQPMessage $message */
  302. $message = $channel->basic_get($queue_name, true);
  303. return $message;
  304. }
  305. /**
  306. * 消息添加延时重试
  307. * @desc 消息消费失败,加入重试队列
  308. * @param object $message
  309. * @param string $exchange_type 交换机类型:默认直连交换机
  310. * @throws Exception
  311. * @return void
  312. */
  313. public function retry($message,$exchange_type = AMQPExchangeType::DIRECT)
  314. {
  315. $message_id = $message->get('message_id');
  316. //headersObject 是一个AMQPTable对象
  317. $headersObject = $message->get_properties()['application_headers'];
  318. $headersArr = $headersObject->getNativeData();
  319. //投递到重试队列
  320. $headersArr['retry'] = intval($headersArr['retry']);
  321. $headersArr['retry']++;
  322. $sec = $headersArr['retry'] * 20; //N个20秒后执行
  323. $micro_sec = $sec * 1000;
  324. if ($headersArr['retry'] <= RabbitMqEnum::MAX_RETRY_NUM) {
  325. //这里加入 延时重试队列
  326. $conn = $this->connect();
  327. $channel = $conn->channel();
  328. $channel->exchange_declare(RabbitMqEnum::EXCHANGE_DELAY_RETRY, $exchange_type, false, false, false);
  329. $channel->exchange_declare(RabbitMqEnum::EXCHANGE_DELAY_RETRY_CACHE, $exchange_type, false, false, false);
  330. //死信交换机和路由
  331. $tale = new AMQPTable();
  332. $tale->set('x-dead-letter-exchange', RabbitMqEnum::EXCHANGE_DELAY_RETRY);
  333. $tale->set('x-dead-letter-routing-key', RabbitMqEnum::DELAY_RETRY_QUEUE);
  334. //$tale->set('x-message-ttl', $micro_sec);
  335. $queue_name = RabbitMqEnum::DELAY_RETRY_QUEUE . '_' . $sec . 's';
  336. //延时缓存队列声明及绑定
  337. $channel->queue_declare($queue_name, false, true, false, true, false, $tale);
  338. $channel->queue_bind($queue_name, RabbitMqEnum::EXCHANGE_DELAY_RETRY_CACHE, $queue_name);
  339. //延时处理队列声明及绑定
  340. $channel->queue_declare(RabbitMqEnum::DELAY_RETRY_QUEUE, false, true, false, false, false);
  341. $channel->queue_bind(RabbitMqEnum::DELAY_RETRY_QUEUE, RabbitMqEnum::EXCHANGE_DELAY_RETRY, RabbitMqEnum::DELAY_RETRY_QUEUE);
  342. $body = new AMQPMessage($message->getBody(), array(
  343. 'expiration' => $micro_sec,
  344. 'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
  345. 'message_id' => $message_id,
  346. 'application_headers' => new AMQPTable($headersArr)
  347. ));
  348. $channel->basic_publish($body, RabbitMqEnum::EXCHANGE_DELAY_RETRY_CACHE, $queue_name);
  349. $channel->close();
  350. $conn->close();
  351. } else {
  352. echo 'retry too many times' . PHP_EOL;
  353. //重试太多次
  354. Yii::error('message retry too many times , message_id :'.$message_id.' ,message content:' . $message->getBody());
  355. Yii::getLogger()->flush(true);
  356. }
  357. }
  358. /**
  359. * 生成MQ消息全局唯一ID
  360. * @return mixed
  361. */
  362. public function createMessageId()
  363. {
  364. $key = RabbitMqEnum::RABBITMQ_QUEUE_ID;
  365. $config = \Yii::$app->getComponents()['redis'];
  366. unset($config['class']);
  367. $redis = new Connection($config);
  368. $id = $redis->incr($key);
  369. $redis->close();
  370. return $id;
  371. }
  372. /**
  373. * redis设置消息信息
  374. * @param $message_id
  375. * @param $data
  376. * @return bool
  377. * @throws
  378. */
  379. public function setMessage($message_id, $data,$expire_time=86400)
  380. {
  381. $key = RabbitMqEnum::QUEUE_MESSAGE_KEY_PREFIX . $message_id;
  382. $config = \Yii::$app->getComponents()['redis'];
  383. unset($config['class']);
  384. $redis = new Connection($config);
  385. $res = $redis->executeCommand('SET', [$key, json_encode($data), 'EX', $expire_time, 'NX']);
  386. $redis->close();
  387. return $res;
  388. }
  389. /**
  390. * 判断消息是否存在(不存在:已消费)
  391. * @param $message_id
  392. * @return mixed
  393. * @throws
  394. */
  395. public function checkMessage($message_id)
  396. {
  397. $key = RabbitMqEnum::QUEUE_MESSAGE_KEY_PREFIX . $message_id;
  398. $config = \Yii::$app->getComponents()['redis'];
  399. unset($config['class']);
  400. $redis = new Connection($config);
  401. $res = $redis->executeCommand('GET', [$key]);
  402. $redis->close();
  403. return $res;
  404. }
  405. /**
  406. * 消息消费完,删除
  407. * @param $message_id
  408. * @return mixed
  409. */
  410. public function deleteMessage($message_id)
  411. {
  412. $key = RabbitMqEnum::QUEUE_MESSAGE_KEY_PREFIX . $message_id;
  413. $config = \Yii::$app->getComponents()['redis'];
  414. unset($config['class']);
  415. $redis = new Connection($config);
  416. $res = $redis->del($key);
  417. $redis->close();
  418. return $res;
  419. }
  420. }