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

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

发布:smiling 来源: PHP粉丝网  添加日期:2021-11-13 22:29:26 浏览: 评论:0 

这篇文章主要给大家介绍了关于利用PHP+RabbitMQ实现消息队列的相关资料,文中通过示例代码介绍的非常详细,对大家学习或者使用PHP具有一定的参考学习价值,需要的朋友们下面来一起学习学习吧。

前言

为什么使用RabbitMq而不是ActiveMq或者RocketMq?

首先,从业务上来讲,我并不要求消息的100%接受率,并且,我需要结合php开发,RabbitMq相较RocketMq,延迟较低(微妙级)。至于ActiveMq,貌似问题较多。RabbitMq对各种语言的支持较好,所以选择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. return [ 
  3.  //配置 
  4.  'host' => [ 
  5.   'host' => '127.0.0.1', 
  6.   'port' => '5672', 
  7.   'login' => 'guest', 
  8.   'password' => 'guest', 
  9.   'vhost'=>'/', 
  10.  ], 
  11.  //交换机 
  12.  'exchange'=>'word', 
  13.  //路由 
  14.  'routes' => [], 
  15. ]; 

BaseMQ.php

  1. <?php 
  2. /** 
  3.  * Created by PhpStorm. 
  4.  * User: pc 
  5.  * Date: 2018/12/13 
  6.  * Time: 14:11 
  7.  */ 
  8.  
  9. namespace MyObjSummary\rabbitMQ; 
  10.  
  11. /** Member 
  12.  *  AMQPChannel 
  13.  *  AMQPConnection 
  14.  *  AMQPEnvelope 
  15.  *  AMQPExchange 
  16.  *  AMQPQueue 
  17.  * Class BaseMQ 
  18.  * @package MyObjSummary\rabbitMQ 
  19.  */ 
  20. class BaseMQ 
  21. { 
  22.  /** MQ Channel 
  23.   * @var \AMQPChannel 
  24.   */ 
  25.  public $AMQPChannel ; 
  26.  
  27.  /** MQ Link 
  28.   * @var \AMQPConnection 
  29.   */ 
  30.  public $AMQPConnection ; 
  31.  
  32.  /** MQ Envelope 
  33.   * @var \AMQPEnvelope 
  34.   */ 
  35.  public $AMQPEnvelope ; 
  36.  
  37.  /** MQ Exchange 
  38.   * @var \AMQPExchange 
  39.   */ 
  40.  public $AMQPExchange ; 
  41.  
  42.  /** MQ Queue 
  43.   * @var \AMQPQueue 
  44.   */ 
  45.  public $AMQPQueue ; 
  46.  
  47.  /** conf 
  48.   * @var 
  49.   */ 
  50.  public $conf ; 
  51.  
  52.  /** exchange 
  53.   * @var 
  54.   */ 
  55.  public $exchange ; 
  56.  
  57.  /** link 
  58.   * BaseMQ constructor. 
  59.   * @throws \AMQPConnectionException 
  60.   */ 
  61.  public function __construct() 
  62.  { 
  63.   $conf = require 'config.php' ; 
  64.   if(!$conf) 
  65.    throw new \AMQPConnectionException('config error!'); 
  66.   $this->conf  = $conf['host'] ; 
  67.   $this->exchange = $conf['exchange'] ; 
  68.   $this->AMQPConnection = new \AMQPConnection($this->conf); 
  69.   if (!$this->AMQPConnection->connect()) 
  70.    throw new \AMQPConnectionException("Cannot connect to the broker!\n"); 
  71.  } 
  72.  
  73.  /** 
  74.   * close link 
  75.   */ 
  76.  public function close() 
  77.  { 
  78.   $this->AMQPConnection->disconnect(); 
  79.  } 
  80.  
  81.  /** Channel 
  82.   * @return \AMQPChannel 
  83.   * @throws \AMQPConnectionException 
  84.   */ 
  85.  public function channel() 
  86.  { 
  87.   if(!$this->AMQPChannel) { 
  88.    $this->AMQPChannel = new \AMQPChannel($this->AMQPConnection); 
  89.   } 
  90.   return $this->AMQPChannel; 
  91.  } 
  92.  
  93.  /** Exchange 
  94.   * @return \AMQPExchange 
  95.   * @throws \AMQPConnectionException 
  96.   * @throws \AMQPExchangeException 
  97.   */ 
  98.  public function exchange() 
  99.  { 
  100.   if(!$this->AMQPExchange) { 
  101.    $this->AMQPExchange = new \AMQPExchange($this->channel()); 
  102.    $this->AMQPExchange->setName($this->exchange); 
  103.   } 
  104.   return $this->AMQPExchange ; 
  105.  } 
  106.  
  107.  /** queue 
  108.   * @return \AMQPQueue 
  109.   * @throws \AMQPConnectionException 
  110.   * @throws \AMQPQueueException 
  111.   */ 
  112.  public function queue() 
  113.  { 
  114.   if(!$this->AMQPQueue) { 
  115.    $this->AMQPQueue = new \AMQPQueue($this->channel()); 
  116.   } 
  117.   return $this->AMQPQueue ; 
  118.  } 
  119.  
  120.  /** Envelope 
  121.   * @return \AMQPEnvelope 
  122.   */ 
  123.  public function envelope() 
  124.  { 
  125.   if(!$this->AMQPEnvelope) { 
  126.    $this->AMQPEnvelope = new \AMQPEnvelope(); 
  127.   } 
  128.   return $this->AMQPEnvelope; 
  129.  } 
  130. } 

ProductMQ.php

  1. <?php 
  2. //生产者 P 
  3. namespace MyObjSummary\rabbitMQ; 
  4. require 'BaseMQ.php'; 
  5. class ProductMQ extends BaseMQ 
  6. { 
  7.  private $routes = ['hello','word']; //路由key 
  8.  
  9.  /** 
  10.   * ProductMQ constructor. 
  11.   * @throws \AMQPConnectionException 
  12.   */ 
  13.  public function __construct() 
  14.  { 
  15.   parent::__construct(); 
  16.  } 
  17.  
  18.  /** 只控制发送成功 不接受消费者是否收到 
  19.   * @throws \AMQPChannelException 
  20.   * @throws \AMQPConnectionException 
  21.   * @throws \AMQPExchangeException 
  22.   */ 
  23.  public function run() 
  24.  { 
  25.   //频道 
  26.   $channel = $this->channel(); 
  27.   //创建交换机对象 
  28.   $ex = $this->exchange(); 
  29.   //消息内容 
  30.   $message = 'product message '.rand(1,99999); 
  31.   //开始事务 
  32.   $channel->startTransaction(); 
  33.   $sendEd = true ; 
  34.   foreach ($this->routes as $route) { 
  35.    $sendEd = $ex->publish($message, $route) ; 
  36.    echo "Send Message:".$sendEd."\n"; 
  37.   } 
  38.   if(!$sendEd) { 
  39.    $channel->rollbackTransaction(); 
  40.   } 
  41.   $channel->commitTransaction(); //提交事务 
  42.   $this->close(); 
  43.   die ; 
  44.  } 
  45. } 
  46. try{ 
  47.  (new ProductMQ())->run(); 
  48. }catch (\Exception $exception){ 
  49.  var_dump($exception->getMessage()) ; 
  50. } 

ConsumerMQ.php

  1. <?php 
  2. //消费者 C 
  3. namespace MyObjSummary\rabbitMQ; 
  4. require 'BaseMQ.php'; 
  5. class ConsumerMQ extends BaseMQ 
  6. { 
  7.  private $q_name = 'hello'; //队列名 
  8.  private $route = 'hello'; //路由key 
  9.  
  10.  /** 
  11.   * ConsumerMQ constructor. 
  12.   * @throws \AMQPConnectionException 
  13.   */ 
  14.  public function __construct() 
  15.  { 
  16.   parent::__construct(); 
  17.  } 
  18.  
  19.  /** 接受消息 如果终止 重连时会有消息 
  20.   * @throws \AMQPChannelException 
  21.   * @throws \AMQPConnectionException 
  22.   * @throws \AMQPExchangeException 
  23.   * @throws \AMQPQueueException 
  24.   */ 
  25.  public function run() 
  26.  { 
  27.  
  28.   //创建交换机 
  29.   $ex = $this->exchange(); 
  30.   $ex->setType(AMQP_EX_TYPE_DIRECT); //direct类型 
  31.   $ex->setFlags(AMQP_DURABLE); //持久化 
  32.   //echo "Exchange Status:".$ex->declare()."\n"; 
  33.  
  34.   //创建队列 
  35.   $q = $this->queue(); 
  36.   //var_dump($q->declare());exit(); 
  37.   $q->setName($this->q_name); 
  38.   $q->setFlags(AMQP_DURABLE); //持久化 
  39.   //echo "Message Total:".$q->declareQueue()."\n"; 
  40.  
  41.   //绑定交换机与队列,并指定路由键 
  42.   echo 'Queue Bind: '.$q->bind($this->exchange, $this->route)."\n"; 
  43.  
  44.   //阻塞模式接收消息 
  45.   echo "Message:\n"; 
  46.   while(True){ 
  47.    $q->consume(function ($envelope,$queue){ 
  48.     $msg = $envelope->getBody(); 
  49.     echo $msg."\n"; //处理消息 
  50.     $queue->ack($envelope->getDeliveryTag()); //手动发送ACK应答 
  51.    }); 
  52.    //$q->consume('processMessage', AMQP_AUTOACK); //自动ACK应答 
  53.   } 
  54.   $this->close(); 
  55.  } 
  56. } 
  57. try{ 
  58.  (new ConsumerMQ)->run(); 
  59. }catch (\Exception $exception){ 
  60.  var_dump($exception->getMessage()) ; 
  61. }

Tags: PHP+RabbitMQ PHP消息队列

分享到: