在PHP中如何使用RabbitMQ来实现消息的订阅和发布?
bigegpt 2024-10-12 06:59 18 浏览
本文将介绍在PHP中如何使用RabbitMQ来实现消息的订阅和发布。我使用的系统依然是Centos7,为了方便,应用服务器我使用Docker进行部署,容器环境:centos7+nginx+php5.6。
运行环境,安装AMQP扩展:
如何安装Docker我就不说了,网上很多教程非常简单,如果有现成的php环境可以直接使用。Docker中我使用的镜像名为webdevops/php-nginx,tag为:centos-7-php56。下载镜像:
(国际带宽出口不稳定,可能会下载失败,重试记次就好了)
docker pull webdevops/php-nginx:centos-7-php56 //下载镜像 docker run -d -p 80:80 --name rabbitmq webdevops/php-nginx:centos-7-php56 //运行容器 docker exec -ti rabbitmq /bin/bash //进入容器
进入到容器后检测下环境是否有相应扩展
cd app vi index.php
刚刚我们在运行容器的时候使用80端口,在浏览器中输入http://ip
搜索下没有amqp相关的信息。下面开始安装amqp扩展。
yum install gcc librabbitmq-devel.x86_64 php56w-devel -y wget http://pecl.php.net/get/amqp-1.4.0.tgz tar -zxvf amqp-1.4.0.tgz cd amqp-1.4.0 phpize ./configure --with-amqp make && make install
在php.ini中开启extension=amqp.so 接着重启php-fpm 或 Web服务器
vi /etc/php.ini extension=amqp.so
我这里就直接重启容器了,如果是宿主机直接安装php环境直接重启环境。
exit //退出容器 docker restart rabbitmq //重启容器
再查看phpinfo,amqp扩展已经安装好了:
publish发布消息
在/app路径下新建一个publish.php的文件
touch publish.php vi publish.php
以下是PHP代码,我们先定义好用来发消息的交换机、队列、RoutingKey、消息等变量。
$queueName = 'superrd'; $exchangeName = 'superrd'; $routeKey = 'superrd'; $message = 'Hello World!';
按照我们第二章讲到的首先建立一个连接。
$connection = new AMQPConnection(array('host' => '10.99.121.137', 'port' => '5672', 'vhost' => '/', 'login' => 'superrd', 'password' => 'superrd')); $connection->connect() or die("Cannot connect to the broker!\n");
新建一个信道。
$channel = new AMQPChannel($connection);
新建一个交换机Exchange,并定义属性,第二章我们讲过有四种类型的交换机,这里使用直连型DIRECT。AMQP_DURABLE代表这是一个持久化的交换机,不会以为服务器异常等因素丢失。
$exchange = new AMQPExchange($channel); $exchange->setName($exchangeName); $exchange->setType(AMQP_EX_TYPE_DIRECT); $exchange->setFlags(AMQP_DURABLE); $exchange->declareExchange();
新建一个队列Queue,前面也讲过生产者将消息发送到Exchange中,Exchange会根据绑定关系投递到队列,也就是如果生产者在生产消息时没有队列与之绑定消息就会丢失。为了保证系统更加健硕,一般无论是消息的生产者还是消费者都会新建一遍Exchange和Queue,新建后属性不会改变。同样AMQP_DURABLE代表这是一个持久化的队列,队列会被写入磁盘。需要注意的是虽然消息是缓存在队列中,但是并不是队列是持久化的队列队列中的消息就是持久化的,消息的持久化需要单独设置。
$queue = new AMQPQueue($channel); $queue->setName($queueName); $queue->setFlags(AMQP_DURABLE); $queue->declareQueue();
通过routeKey绑定交换机和队列。
$queue->bind($exchangeName, $routeKey);
好了,下面可以发送消息了
$exchange->publish($message,$routeKey);
如果你希望消息也是持久化的可以使用如下的代码,实际测试结果在持久化消息后消息发布的性能下降一倍,我的磁盘是pcie的固态硬盘,如果你是机械磁盘这个性能下降估计会更明显,24核心CPU,48GB内存,pcie固态硬盘,单线程的情况下每秒可以发布2.5万左右的非持久化消息,持久化之后变为变为1.2万左右。
$exchange->publish($message,$routeKey,AMQP_NOPARAM, array('delivery_mode'=>2));
断开连接。
$connection->disconnect();
同样在发布消息之后可以通过WEB工具来查看是否发布成功,
查看交换机多了一个superid交换机。
查看交换机已经有superrd队列。
点击队列查看队列详情。Bindings标签可以看到交换机和队列的绑定关系。
点击Get messages标签Get message(s)按钮可以看到队列中的消息。
到此说明我们已经将一个消息发布到了消息队列中。完整的PHP代码如下。
'10.99.121.137', 'port' => '5672', 'vhost' => '/', 'login' => 'superrd', 'password' => 'superrd')); $connection->connect() or die("Cannot connect to the broker!\n"); try { $channel = new AMQPChannel($connection); $exchange = new AMQPExchange($channel); $exchange->setName($exchangeName); $exchange->setType(AMQP_EX_TYPE_DIRECT); $exchange->setFlags(AMQP_DURABLE); $exchange->declareExchange(); $queue = new AMQPQueue($channel); $queue->setName($queueName); $queue->setFlags(AMQP_DURABLE); $queue->declareQueue(); $queue->bind($exchangeName, $routeKey); $exchange->publish($message,$routeKey); var_dump("[x] Sent 'Hello World!'"); } catch (AMQPConnectionException $e) { var_dump($e); exit(); } $connection->disconnect();
Subscribe订阅消息
在/app路径下新建一个subscribe.php的文件
touch subscribe.php vi subscribe.php
以下是PHP代码,和发布消息一样我们先定义好用交换机、队列、RoutingKey等变量。
$queueName = 'superrd'; $exchangeName = 'superrd'; $routeKey = 'superrd';
按照我们第二章讲到的首先建立一个连接。
$connection = new AMQPConnection(array('host' => '10.99.121.137', 'port' => '5672', 'vhost' => '/', 'login' => 'superrd', 'password' => 'superrd')); $connection->connect() or die("Cannot connect to the broker!\n");
新建一个信道。
$channel = new AMQPChannel($connection);
与发布消息一样新建交换机。
$exchange = new AMQPExchange($channel); $exchange->setName($exchangeName); $exchange->setType(AMQP_EX_TYPE_DIRECT); $exchange->setFlags(AMQP_DURABLE); $exchange->declareExchange();
新建一个队列Queue。
$queue = new AMQPQueue($channel); $queue->setName($queueName); $queue->setFlags(AMQP_DURABLE); $queue->declareQueue();
通过routeKey绑定交换机和队列。
$queue->bind($exchangeName, $routeKey);
重点来了,阻塞订阅消息。
//阻塞模式接收消息 echo "Message:\n"; while(True){ $queue->consume('processMessage'); //自动ACK应答 //$queue->consume('processMessage', AMQP_AUTOACK); } $conn->disconnect(); /* * 消费回调函数 * 处理消息 */ function processMessage($envelope, $q) { $msg = $envelope->getBody(); echo $msg."\n"; //处理消息 $q->ack($envelope->getDeliveryTag()); //手动发送ACK应答 }
注意因为是阻塞监听,因为输出缓冲区的原因用浏览器访问该文件是看不到输出的。使用脚本访问。
php /app/subscribe.php
通过WEB工具查看队列。superrd队列中的消息数已经为0。
完整的PHP代码如下。
'10.99.121.137', 'port' => '5672', 'vhost' => '/', 'login' => 'superrd', 'password' => 'superrd')); $connection->connect() or die("Cannot connect to the broker!\n"); $channel = new AMQPChannel($connection); $exchange = new AMQPExchange($channel); $exchange->setName($exchangeName); $exchange->setType(AMQP_EX_TYPE_DIRECT); $exchange->setFlags(AMQP_DURABLE); $exchange->declareExchange(); $queue = new AMQPQueue($channel); $queue->setName($queueName); $queue->setFlags(AMQP_DURABLE); $queue->declareQueue(); $queue->bind($exchangeName, $routeKey); //阻塞模式接收消息 echo "Message:\n"; while(True){ $queue->consume('processMessage'); //自动ACK应答 //$queue->consume('processMessage', AMQP_AUTOACK); } $conn->disconnect(); /* * 消费回调函数 * 处理消息 */ function processMessage($envelope, $q) { $msg = $envelope->getBody(); echo $msg."\n"; //处理消息 $q->ack($envelope->getDeliveryTag()); //手动发送ACK应答 }
相关推荐
- 当Frida来“敲”门(frida是什么)
-
0x1渗透测试瓶颈目前,碰到越来越多的大客户都会将核心资产业务集中在统一的APP上,或者对自己比较重要的APP,如自己的主业务,办公APP进行加壳,流量加密,投入了很多精力在移动端的防护上。而现在挖...
- 服务端性能测试实战3-性能测试脚本开发
-
前言在前面的两篇文章中,我们分别介绍了性能测试的理论知识以及性能测试计划制定,本篇文章将重点介绍性能测试脚本开发。脚本开发将分为两个阶段:阶段一:了解各个接口的入参、出参,使用Python代码模拟前端...
- Springboot整合Apache Ftpserver拓展功能及业务讲解(三)
-
今日分享每天分享技术实战干货,技术在于积累和收藏,希望可以帮助到您,同时也希望获得您的支持和关注。架构开源地址:https://gitee.com/msxyspringboot整合Ftpserver参...
- Linux和Windows下:Python Crypto模块安装方式区别
-
一、Linux环境下:fromCrypto.SignatureimportPKCS1_v1_5如果导包报错:ImportError:Nomodulenamed'Crypt...
- Python 3 加密简介(python des加密解密)
-
Python3的标准库中是没多少用来解决加密的,不过却有用于处理哈希的库。在这里我们会对其进行一个简单的介绍,但重点会放在两个第三方的软件包:PyCrypto和cryptography上,我...
- 怎样从零开始编译一个魔兽世界开源服务端Windows
-
第二章:编译和安装我是艾西,上期我们讲述到编译一个魔兽世界开源服务端环境准备,那么今天跟大家聊聊怎么编译和安装我们直接进入正题(上一章没有看到的小伙伴可以点我主页查看)编译服务端:在D盘新建一个文件夹...
- 附1-Conda部署安装及基本使用(conda安装教程)
-
Windows环境安装安装介质下载下载地址:https://www.anaconda.com/products/individual安装Anaconda安装时,选择自定义安装,选择自定义安装路径:配置...
- 如何配置全世界最小的 MySQL 服务器
-
配置全世界最小的MySQL服务器——如何在一块IntelEdison为控制板上安装一个MySQL服务器。介绍在我最近的一篇博文中,物联网,消息以及MySQL,我展示了如果Partic...
- 如何使用Github Action来自动化编译PolarDB-PG数据库
-
随着PolarDB在国产数据库领域荣膺桂冠并持续获得广泛认可,越来越多的学生和技术爱好者开始关注并涉足这款由阿里巴巴集团倾力打造且性能卓越的关系型云原生数据库。有很多同学想要上手尝试,却卡在了编译数据...
- 面向NDK开发者的Android 7.0变更(ndk android.mk)
-
订阅Google官方微信公众号:谷歌开发者。与谷歌一起创造未来!受Android平台其他改进的影响,为了方便加载本机代码,AndroidM和N中的动态链接器对编写整洁且跨平台兼容的本机...
- 信创改造--人大金仓(Kingbase)数据库安装、备份恢复的问题纪要
-
问题一:在安装KingbaseES时,安装用户对于安装路径需有“读”、“写”、“执行”的权限。在Linux系统中,需要以非root用户执行安装程序,且该用户要有标准的home目录,您可...
- OpenSSH 安全漏洞,修补操作一手掌握
-
1.漏洞概述近日,国家信息安全漏洞库(CNNVD)收到关于OpenSSH安全漏洞(CNNVD-202407-017、CVE-2024-6387)情况的报送。攻击者可以利用该漏洞在无需认证的情况下,通...
- Linux:lsof命令详解(linux lsof命令详解)
-
介绍欢迎来到这篇博客。在这篇博客中,我们将学习Unix/Linux系统上的lsof命令行工具。命令行工具是您使用CLI(命令行界面)而不是GUI(图形用户界面)运行的程序或工具。lsoflsof代表&...
- 幻隐说固态第一期:固态硬盘接口类别
-
前排声明所有信息来源于网络收集,如有错误请评论区指出更正。废话不多说,目前固态硬盘接口按速度由慢到快分有这几类:SATA、mSATA、SATAExpress、PCI-E、m.2、u.2。下面我们来...
- 新品轰炸 影驰SSD多款产品登Computex
-
分享泡泡网SSD固态硬盘频道6月6日台北电脑展作为全球第二、亚洲最大的3C/IT产业链专业展,吸引了众多IT厂商和全球各地媒体的热烈关注,全球存储新势力—影驰,也积极参与其中,为广大玩家朋友带来了...
- 一周热门
- 最近发表
-
- 当Frida来“敲”门(frida是什么)
- 服务端性能测试实战3-性能测试脚本开发
- Springboot整合Apache Ftpserver拓展功能及业务讲解(三)
- Linux和Windows下:Python Crypto模块安装方式区别
- Python 3 加密简介(python des加密解密)
- 怎样从零开始编译一个魔兽世界开源服务端Windows
- 附1-Conda部署安装及基本使用(conda安装教程)
- 如何配置全世界最小的 MySQL 服务器
- 如何使用Github Action来自动化编译PolarDB-PG数据库
- 面向NDK开发者的Android 7.0变更(ndk android.mk)
- 标签列表
-
- mybatiscollection (79)
- mqtt服务器 (88)
- keyerror (78)
- c#map (65)
- resize函数 (64)
- xftp6 (83)
- bt搜索 (75)
- c#var (76)
- mybatis大于等于 (64)
- xcode-select (66)
- mysql授权 (74)
- 下载测试 (70)
- linuxlink (65)
- pythonwget (67)
- androidinclude (65)
- libcrypto.so (74)
- logstashinput (65)
- hadoop端口 (65)
- vue阻止冒泡 (67)
- jquery跨域 (68)
- php写入文件 (73)
- kafkatools (66)
- mysql导出数据库 (66)
- jquery鼠标移入移出 (71)
- 取小数点后两位的函数 (73)