国产探花免费观看_亚洲丰满少妇自慰呻吟_97日韩有码在线_资源在线日韩欧美_一区二区精品毛片,辰东完美世界有声小说,欢乐颂第一季,yy玄幻小说排行榜完本

首頁 > 語言 > PHP > 正文

PHP基于rabbitmq操作類的生產者和消費者功能示例

2024-05-05 00:04:16
字體:
來源:轉載
供稿:網友

本文實例講述了PHP基于rabbitmq操作類的生產者和消費者功能。分享給大家供大家參考,具體如下:

注意事項:

1、accept.php消費者代碼需要在命令行執(zhí)行

2、'username'=>'asdf','password'=>'123456' 改成自己的帳號和密碼

RabbitMQCommand.php操作類代碼

<?php/* * amqp協(xié)議操作類,可以訪問rabbitMQ * 需先安裝php_amqp擴展 */class RabbitMQCommand{  public $configs = array();  //交換機名稱  public $exchange_name = '';  //隊列名稱  public $queue_name = '';  //路由名稱  public $route_key = '';  /*   * 持久化,默認True   */  public $durable = True;  /*   * 自動刪除   * exchange is deleted when all queues have finished using it   * queue is deleted when last consumer unsubscribes   *   */  public $autodelete = False;  /*   * 鏡像   * 鏡像隊列,打開后消息會在節(jié)點之間復制,有master和slave的概念   */  public $mirror = False;  private $_conn = Null;  private $_exchange = Null;  private $_channel = Null;  private $_queue = Null;  /*   * @configs array('host'=>$host,'port'=>5672,'username'=>$username,'password'=>$password,'vhost'=>'/')   */  public function __construct($configs = array(), $exchange_name = '', $queue_name = '', $route_key = '') {    $this->setConfigs($configs);    $this->exchange_name = $exchange_name;    $this->queue_name = $queue_name;    $this->route_key = $route_key;  }  private function setConfigs($configs) {    if (!is_array($configs)) {      throw new Exception('configs is not array');    }    if (!($configs['host'] && $configs['port'] && $configs['username'] && $configs['password'])) {      throw new Exception('configs is empty');    }    if (empty($configs['vhost'])) {      $configs['vhost'] = '/';    }    $configs['login'] = $configs['username'];    unset($configs['username']);    $this->configs = $configs;  }  /*   * 設置是否持久化,默認為True   */  public function setDurable($durable) {    $this->durable = $durable;  }  /*   * 設置是否自動刪除   */  public function setAutoDelete($autodelete) {    $this->autodelete = $autodelete;  }  /*   * 設置是否鏡像   */  public function setMirror($mirror) {    $this->mirror = $mirror;  }  /*   * 打開amqp連接   */  private function open() {    if (!$this->_conn) {      try {        $this->_conn = new AMQPConnection($this->configs);        $this->_conn->connect();        $this->initConnection();      } catch (AMQPConnectionException $ex) {        throw new Exception('cannot connection rabbitmq',500);      }    }  }  /*   * rabbitmq連接不變   * 重置交換機,隊列,路由等配置   */  public function reset($exchange_name, $queue_name, $route_key) {    $this->exchange_name = $exchange_name;    $this->queue_name = $queue_name;    $this->route_key = $route_key;    $this->initConnection();  }  /*   * 初始化rabbit連接的相關配置   */  private function initConnection() {    if (empty($this->exchange_name) || empty($this->queue_name) || empty($this->route_key)) {      throw new Exception('rabbitmq exchange_name or queue_name or route_key is empty',500);    }    $this->_channel = new AMQPChannel($this->_conn);    $this->_exchange = new AMQPExchange($this->_channel);    $this->_exchange->setName($this->exchange_name);    $this->_exchange->setType(AMQP_EX_TYPE_DIRECT);    if ($this->durable)      $this->_exchange->setFlags(AMQP_DURABLE);    if ($this->autodelete)      $this->_exchange->setFlags(AMQP_AUTODELETE);    $this->_exchange->declare();    $this->_queue = new AMQPQueue($this->_channel);    $this->_queue->setName($this->queue_name);    if ($this->durable)      $this->_queue->setFlags(AMQP_DURABLE);    if ($this->autodelete)      $this->_queue->setFlags(AMQP_AUTODELETE);    if ($this->mirror)      $this->_queue->setArgument('x-ha-policy', 'all');    $this->_queue->declare();    $this->_queue->bind($this->exchange_name, $this->route_key);  }  public function close() {    if ($this->_conn) {      $this->_conn->disconnect();    }  }  public function __sleep() {    $this->close();    return array_keys(get_object_vars($this));  }  public function __destruct() {    $this->close();  }  /*   * 生產者發(fā)送消息   */  public function send($msg) {    $this->open();    if(is_array($msg)){      $msg = json_encode($msg);    }else{      $msg = trim(strval($msg));    }    return $this->_exchange->publish($msg, $this->route_key);  }  /*   * 消費者   * $fun_name = array($classobj,$function) or function name string   * $autoack 是否自動應答   *   * function processMessage($envelope, $queue) {      $msg = $envelope->getBody();      echo $msg."/n"; //處理消息      $queue->ack($envelope->getDeliveryTag());//手動應答    }   */  public function run($fun_name, $autoack = True){    $this->open();    if (!$fun_name || !$this->_queue) return False;    while(True){      if ($autoack) $this->_queue->consume($fun_name, AMQP_AUTOACK);      else $this->_queue->consume($fun_name);    }  }}

send.php生產者代碼

<?phpset_time_limit(0);include_once('RabbitMQCommand.php');$configs = array('host'=>'127.0.0.1','port'=>5672,'username'=>'asdf','password'=>'123456','vhost'=>'/');$exchange_name = 'class-e-1';$queue_name = 'class-q-1';$route_key = 'class-r-1';$ra = new RabbitMQCommand($configs,$exchange_name,$queue_name,$route_key);for($i=0;$i<=100;$i++){  $ra->send(date('Y-m-d H:i:s',time()));}exit();

accept.php消費者代碼

<?phperror_reporting(0);include_once('RabbitMQCommand.php');$configs = array('host'=>'127.0.0.1','port'=>5672,'username'=>'asdf','password'=>'123456','vhost'=>'/');$exchange_name = 'class-e-1';$queue_name = 'class-q-1';$route_key = 'class-r-1';$ra = new RabbitMQCommand($configs,$exchange_name,$queue_name,$route_key);class A{  function processMessage($envelope, $queue) {    $msg = $envelope->getBody();    $envelopeID = $envelope->getDeliveryTag();    $pid = posix_getpid();    file_put_contents("log{$pid}.log", $msg.'|'.$envelopeID.''."/r/n",FILE_APPEND);    $queue->ack($envelopeID);  }}$a = new A();$s = $ra->run(array($a,'processMessage'),false);

希望本文所述對大家PHP程序設計有所幫助。


注:相關教程知識閱讀請移步到PHP教程頻道。
發(fā)表評論 共有條評論
用戶名: 密碼:
驗證碼: 匿名發(fā)表

圖片精選

主站蜘蛛池模板: 贺兰县| 晋中市| 凤山市| 连江县| 筠连县| 泾川县| 丘北县| 三原县| 昆明市| 揭阳市| 航空| 宜丰县| 九龙城区| 泗水县| 司法| 徐水县| 青岛市| 信宜市| 贵德县| 富阳市| 武夷山市| 浏阳市| 盱眙县| 新绛县| 攀枝花市| 蕲春县| 伽师县| 武平县| 阜康市| 青海省| 怀远县| 龙山县| 夹江县| 大洼县| 遂溪县| 新巴尔虎右旗| 板桥市| 察雅县| 淮安市| 陆丰市| 栾川县|