当前位置:首页 > PHP教程 > php高级应用 > 列表

PHP和RabbitMQ实现消息队列的完整代码

发布:smiling 来源: PHP粉丝网  添加日期:2020-02-04 16:27:05 浏览: 评论:0 

本篇文章给大家带来的内容是关于PHP和RabbitMQ实现消息队列的完整代码,有一定的参考价值,有需要的朋友可以参考一下,希望对你有所帮助。

先安装PHP对应的RabbitMQ,这里用的是 php_amqp 不同的扩展实现方式会有细微的差异.

php扩展地址: http://pecl.php.net/package/amqp

具体以官网为准 http://www.rabbitmq.com/getstarted.html

介绍:

config.php 配置信息

BaseMQ.php MQ基类

ProductMQ.php 生产者类

ConsumerMQ.php 消费者类

Consumer2MQ.php 消费者2(可有多个)

config.php

  1. <?php 
  2.  
  3. return [ 
  4.  
  5.     //配置 
  6.  
  7.     'host' => [ 
  8.  
  9.         'host' => '127.0.0.1', 
  10.  
  11.         'port' => '5672', 
  12.  
  13.         'login' => 'guest', 
  14.  
  15.         'password' => 'guest', 
  16.  
  17.         'vhost'=>'/', 
  18.  
  19.     ], 
  20.  
  21.     //交换机 
  22. //phpfensi.com 
  23.     'exchange'=>'word', 
  24.  
  25.     //路由 
  26.  
  27.     'routes' => [], 
  28.  
  29. ]; 

BaseMQ.php

  1. <?php 
  2.  
  3. /** 
  4.  
  5.  * Created by PhpStorm. 
  6.  
  7.  * User: pc 
  8.  
  9.  * Date: 2018/12/13 
  10.  
  11.  * Time: 14:11 
  12.  
  13.  */ 
  14.  
  15.  
  16.  
  17. namespace MyObjSummary\rabbitMQ; 
  18.  
  19.  
  20.  
  21. /** Member 
  22.  
  23.  *      AMQPChannel 
  24.  
  25.  *      AMQPConnection 
  26.  
  27.  *      AMQPEnvelope 
  28.  
  29.  *      AMQPExchange 
  30.  
  31.  *      AMQPQueue 
  32.  
  33.  * Class BaseMQ 
  34.  
  35.  * @package MyObjSummary\rabbitMQ 
  36.  
  37.  */ 
  38.  
  39. class BaseMQ 
  40.  
  41. { 
  42.  
  43.     /** MQ Channel 
  44.  
  45.      * @var \AMQPChannel 
  46.  
  47.      */ 
  48.  
  49.     public $AMQPChannel ; 
  50.  
  51.  
  52.  
  53.     /** MQ Link 
  54.  
  55.      * @var \AMQPConnection 
  56.  
  57.      */ 
  58.  
  59.     public $AMQPConnection ; 
  60.  
  61.  
  62.  
  63.     /** MQ Envelope 
  64.  
  65.      * @var \AMQPEnvelope 
  66.  
  67.      */ 
  68.  
  69.     public $AMQPEnvelope ; 
  70.  
  71.  
  72.  
  73.     /** MQ Exchange 
  74.  
  75.      * @var \AMQPExchange 
  76.  
  77.      */ 
  78.  
  79.     public $AMQPExchange ; 
  80.  
  81.  
  82.  
  83.     /** MQ Queue 
  84.  
  85.      * @var \AMQPQueue 
  86.  
  87.      */ 
  88.  
  89.     public $AMQPQueue ; 
  90.  
  91.  
  92.  
  93.     /** conf 
  94.  
  95.      * @var 
  96.  
  97.      */ 
  98.  
  99.     public $conf ; 
  100.  
  101.  
  102.  
  103.     /** exchange 
  104.  
  105.      * @var 
  106.  
  107.      */ 
  108.  
  109.     public $exchange ; 
  110.  
  111.  
  112.  
  113.     /** link 
  114.  
  115.      * BaseMQ constructor. 
  116.  
  117.      * @throws \AMQPConnectionException 
  118.  
  119.      */ 
  120.  
  121.     public function __construct() 
  122.  
  123.     { 
  124.  
  125.         $conf =  require 'config.php' ; 
  126.  
  127.         if(!$conf) 
  128.  
  129.             throw new \AMQPConnectionException('config error!'); 
  130.  
  131.         $this->conf     = $conf['host'] ; 
  132.  
  133.         $this->exchange = $conf['exchange'] ; 
  134.  
  135.         $this->AMQPConnection = new \AMQPConnection($this->conf); 
  136.  
  137.         if (!$this->AMQPConnection->connect()) 
  138.  
  139.             throw new \AMQPConnectionException("Cannot connect to the broker!\n"); 
  140.  
  141.     } 
  142.  
  143.  
  144.  
  145.     /** 
  146.  
  147.      * close link 
  148.  
  149.      */ 
  150.  
  151.     public function close() 
  152.  
  153.     { 
  154.  
  155.         $this->AMQPConnection->disconnect(); 
  156.  
  157.     } 
  158.  
  159.  
  160.  
  161.     /** Channel 
  162.  
  163.      * @return \AMQPChannel 
  164.  
  165.      * @throws \AMQPConnectionException 
  166.  
  167.      */ 
  168.  
  169.     public function channel() 
  170.  
  171.     { 
  172.  
  173.         if(!$this->AMQPChannel) { 
  174.  
  175.             $this->AMQPChannel =  new \AMQPChannel($this->AMQPConnection); 
  176.  
  177.         } 
  178.  
  179.         return $this->AMQPChannel; 
  180.  
  181.     } 
  182.  
  183.  
  184.  
  185.     /** Exchange 
  186.  
  187.      * @return \AMQPExchange 
  188.  
  189.      * @throws \AMQPConnectionException 
  190.  
  191.      * @throws \AMQPExchangeException 
  192.  
  193.      */ 
  194.  
  195.     public function exchange() 
  196.  
  197.     { 
  198.  
  199.         if(!$this->AMQPExchange) { 
  200.  
  201.             $this->AMQPExchange = new \AMQPExchange($this->channel()); 
  202.  
  203.             $this->AMQPExchange->setName($this->exchange); 
  204.  
  205.         } 
  206.  
  207.         return $this->AMQPExchange ; 
  208.  
  209.     } 
  210.  
  211.  
  212.  
  213.     /** queue 
  214.  
  215.      * @return \AMQPQueue 
  216.  
  217.      * @throws \AMQPConnectionException 
  218.  
  219.      * @throws \AMQPQueueException 
  220.  
  221.      */ 
  222.  
  223.     public function queue() 
  224.  
  225.     { 
  226.  
  227.         if(!$this->AMQPQueue) { 
  228.  
  229.             $this->AMQPQueue = new \AMQPQueue($this->channel()); 
  230.  
  231.         } 
  232.  
  233.         return $this->AMQPQueue ; 
  234.  
  235.     } 
  236.  
  237.  
  238.  
  239.     /** Envelope 
  240.  
  241.      * @return \AMQPEnvelope 
  242.  
  243.      */ 
  244.  
  245.     public function envelope() 
  246.  
  247.     { 
  248.  
  249.         if(!$this->AMQPEnvelope) { 
  250.  
  251.             $this->AMQPEnvelope = new \AMQPEnvelope(); 
  252. //phpfensi.com 
  253.         } 
  254.  
  255.         return $this->AMQPEnvelope; 
  256.  
  257.     } 
  258.  
  259. } 

ProductMQ.php

  1. <?php 
  2.  
  3. //生产者 P 
  4.  
  5. namespace MyObjSummary\rabbitMQ; 
  6.  
  7. require 'BaseMQ.php'; 
  8.  
  9. class ProductMQ extends BaseMQ 
  10.  
  11. { 
  12.  
  13.     private $routes = ['hello','word']; //路由key 
  14.  
  15.  
  16.  
  17.     /** 
  18.  
  19.      * ProductMQ constructor. 
  20.  
  21.      * @throws \AMQPConnectionException 
  22.  
  23.      */ 
  24.  
  25.     public function __construct() 
  26.  
  27.     { 
  28.  
  29.        parent::__construct(); 
  30.  
  31.     } 
  32.  
  33.  
  34.  
  35.     /** 只控制发送成功 不接受消费者是否收到 
  36.  
  37.      * @throws \AMQPChannelException 
  38.  
  39.      * @throws \AMQPConnectionException 
  40.  
  41.      * @throws \AMQPExchangeException 
  42.  
  43.      */ 
  44.  
  45.     public function run() 
  46.  
  47.     { 
  48.  
  49.         //频道 
  50.  
  51.         $channel = $this->channel(); 
  52.  
  53.         //创建交换机对象 
  54.  
  55.         $ex = $this->exchange(); 
  56.  
  57.         //消息内容 
  58.  
  59.         $message = 'product message '.rand(1,99999); 
  60.  
  61.         //开始事务 
  62.  
  63.         $channel->startTransaction(); 
  64.  
  65.         $sendEd = true ; 
  66.  
  67.         foreach ($this->routes as $route) { 
  68.  
  69.             $sendEd = $ex->publish($message, $route) ; 
  70.  
  71.             echo "Send Message:".$sendEd."\n"; 
  72.  
  73.         } 
  74.  
  75.         if(!$sendEd) { 
  76.  
  77.             $channel->rollbackTransaction(); 
  78.  
  79.         } 
  80.  
  81.         $channel->commitTransaction(); //提交事务 
  82.  
  83.         $this->close(); 
  84.  
  85.         die ; 
  86.  
  87.     } 
  88.  
  89. } 
  90.  
  91. try{ 
  92.  
  93.     (new ProductMQ())->run(); 
  94.  
  95. }catch (\Exception $exception){ 
  96.  
  97.     var_dump($exception->getMessage()) ; 
  98.  
  99. } 

ConsumerMQ.php

  1. <?php 
  2.  
  3. //消费者 C 
  4.  
  5. namespace MyObjSummary\rabbitMQ; 
  6.  
  7. require 'BaseMQ.php'; 
  8.  
  9. class ConsumerMQ extends BaseMQ 
  10.  
  11. { 
  12.  
  13.     private  $q_name = 'hello'; //队列名 
  14.  
  15.     private  $route  = 'hello'; //路由key 
  16.  
  17.  
  18.  
  19.     /** 
  20.  
  21.      * ConsumerMQ constructor. 
  22.  
  23.      * @throws \AMQPConnectionException 
  24.  
  25.      */ 
  26.  
  27.     public function __construct() 
  28.  
  29.     { 
  30.  
  31.         parent::__construct(); 
  32.  
  33.     } 
  34.  
  35.  
  36.  
  37.     /** 接受消息 如果终止 重连时会有消息 
  38.  
  39.      * @throws \AMQPChannelException 
  40.  
  41.      * @throws \AMQPConnectionException 
  42.  
  43.      * @throws \AMQPExchangeException 
  44.  
  45.      * @throws \AMQPQueueException 
  46.  
  47.      */ 
  48.  
  49.     public function run() 
  50.  
  51.     { 
  52.  
  53.  
  54.  
  55.         //创建交换机 
  56.  
  57.         $ex = $this->exchange(); 
  58.  
  59.         $ex->setType(AMQP_EX_TYPE_DIRECT); //direct类型 
  60.  
  61.         $ex->setFlags(AMQP_DURABLE); //持久化 
  62.  
  63.         //echo "Exchange Status:".$ex->declare()."\n"; 
  64.  
  65.  
  66.  
  67.         //创建队列 
  68.  
  69.         $q = $this->queue(); 
  70.  
  71.         //var_dump($q->declare());exit(); 
  72.  
  73.         $q->setName($this->q_name); 
  74.  
  75.         $q->setFlags(AMQP_DURABLE); //持久化 
  76.  
  77.         //echo "Message Total:".$q->declareQueue()."\n"; 
  78.  
  79.  
  80.  
  81.         //绑定交换机与队列,并指定路由键 
  82.  
  83.         echo 'Queue Bind: '.$q->bind($this->exchange, $this->route)."\n"; 
  84.  
  85.  
  86.  
  87.         //阻塞模式接收消息 
  88.  
  89.         echo "Message:\n"; 
  90.  
  91.         while(True){ 
  92.  
  93.             $q->consume(function ($envelope,$queue){ 
  94.  
  95.                 $msg = $envelope->getBody(); 
  96.  
  97.                 echo $msg."\n"; //处理消息 
  98.  
  99.                 $queue->ack($envelope->getDeliveryTag()); //手动发送ACK应答 
  100.  
  101.             }); 
  102.  
  103.             //$q->consume('processMessage', AMQP_AUTOACK); //自动ACK应答 
  104.  
  105.         } 
  106.  
  107.         $this->close(); 
  108.  
  109.     } 
  110.  
  111. } 
  112.  
  113. try{ 
  114.  
  115.     (new ConsumerMQ)->run(); 
  116.  
  117. }catch (\Exception $exception){ 
  118.  
  119.     var_dump($exception->getMessage()) ; 
  120.  
  121. } 

Tags: RabbitMQ PHP消息队列

分享到: