资讯动态

Hyperf AMQP 组件实战指南:基于 RabbitMQ 的生产消费、延迟队列与 RPC 调用

发布时间:2026/10/9 1:39:29 来源:尧图企业网站定制
后端微服务【免费下载链接】hyperf A coroutine framework that focuses on hyperspeed and flexibility. Building microservice or middleware with ease.项目地址https://gitcode.com/gh_mirrors/hy/hyperf点击查看免费下载Hyperf 的 AMQP 组件hyperf/amqp是 AMQP 0-9-1 标准的实现主要用于对接 RabbitMQ 消息中间件为微服务场景提供消息投递、异步消费、延迟队列与基于消息的 RPC 远程过程调用能力。读完本文你将掌握该组件的完整配置项含义、Producer/Consumer 的注解式开发方式、消费并发与 QOS 调优技巧以及如何基于延迟插件实现定时延迟消息、如何基于消息队列实现跨服务 RPC 调用。组件简介与安装hyperf/amqp在 composer.json 中声明依赖php-amqplib/php-amqplib等库通过php-amqplib与 RabbitMQ 建立连接同时利用 Hyperf 的协程能力Swoole/Swow见 src/amqp/src/IO 下的SwooleIO与SwowIO实现了非阻塞的消息收发。使用 Composer 安装composer require hyperf/amqp安装完成后组件会通过 ConfigProvider.php 自动发布默认配置到config/autoload/amqp.php。默认配置详解组件默认配置的结构如下表所示默认值以 src/amqp/publish/amqp.php 发布文件为准配置类型默认值备注enablebooltrue是否启用 AMQP 组件default.hoststringlocalhostRabbitMQ Host可用AMQP_HOST环境变量覆盖default.portint5672端口号default.userstringguest用户名default.passwordstringguest密码default.vhoststring/vhostdefault.open_sslboolfalse是否启用 SSL 连接default.concurrent.limitint2同时消费的最大协程数量default.poolobject—连接池配置default.pool.connectionsint2进程内保持的连接数default.ioclass-stringIOFactory::classIO 工厂Swoole/Swowdefault.paramsobject—底层连接基本配置其中params的关键参数及其含义如下params [ insist false, // 连接失败时是否坚持重试 login_method AMQPLAIN, // 登录认证方式 login_response null, // 登录响应 locale en_US, // 区域设置 connection_timeout 3, // 连接超时秒 // 尽量保持为 heartbeat 数值的两倍 read_write_timeout 6, // 读写超时秒 context null, // 自定义 stream context keepalive true, // 是否开启 TCP keepalive // 尽量保证每个消息的消费时间小于心跳时间 heartbeat 3, // 心跳间隔秒 channel_rpc_timeout 0.0, // Channel RPC 超时 close_on_destruct false, // 析构时是否关闭连接 max_idle_channels 10, // 最大空闲 Channel 数 connection_name null, // 连接名称 ],一份完整的配置示例支持通过环境变量覆盖连接参数?php use Hyperf\Amqp\IO\IOFactory; use function Hyperf\Support\env; return [ enable true, default [ host env(AMQP_HOST, localhost), port (int) env(AMQP_PORT, 5672), user env(AMQP_USER, guest), password env(AMQP_PASSWORD, guest), vhost env(AMQP_VHOST, /), open_ssl false, concurrent [ limit 2, ], pool [ connections 2, ], io IOFactory::class, params [ insist false, login_method AMQPLAIN, login_response null, locale en_US, connection_timeout 3, read_write_timeout 6, context null, keepalive true, heartbeat 3, channel_rpc_timeout 0.0, close_on_destruct false, max_idle_channels 10, connection_name null, ], ], // 可配置第二个连接池 pool2 [ // ... 与 default 相同结构 ], ];从源码结构看concurrent.limit会在 Consumer.php 的getConcurrent()方法中被读取当该值大于 1 时框架会创建Hyperf\Coroutine\Concurrent实例将每条消息的消费逻辑放入并发协程中执行从而限制单个消费者进程内并行消费的最大数量。配置了多个连接池如default与pool2之后可以在producer或consumer的__construct构造函数中通过给$this-poolName赋值来指定使用哪一个连接池。投递消息Producer使用gen:amqp-producer命令快速生成一个生产者php bin/hyperf.php gen:amqp-producer DemoProducer生成的DemoProducer继承Hyperf\Amqp\Message\ProducerMessage通过修改#[Producer]注解的字段可以替换对应的exchange和routingKey。其中payload就是最终投递到消息队列中的数据因此可以随意改写__construct方法只要最后给payload赋值即可。使用#[Producer]注解时需引入use Hyperf\Amqp\Annotation\Producer;命名空间。?php declare(strict_types1); namespace App\Amqp\Producers; use Hyperf\Amqp\Annotation\Producer; use Hyperf\Amqp\Message\ProducerMessage; use App\Models\User; #[Producer(exchange: hyperf, routingKey: hyperf)] class DemoProducer extends ProducerMessage { public function __construct($id) { // 设置不同 pool $this-poolName pool2; $user User::where(id, $id)-first(); $this-payload [ id $id, data $user-toArray() ]; } }从 Annotation/Producer.php 源码可以看到#[Producer]注解支持三个参数exchange、routingKey和pool用于覆盖默认连接池。在 Producer.php 的injectMessageProperty()方法中框架会通过AnnotationCollector读取注解配置并注入到消息对象中——这也是为什么只需在注解中声明 exchange/routingKey而不必在类属性中重复定义。通过 DI 容器获取Hyperf\Amqp\Producer实例即可投递消息?php use Hyperf\Amqp\Producer; use App\Amqp\Producers\DemoProducer; use Hyperf\Context\ApplicationContext; $message new DemoProducer(1); $producer ApplicationContext::getContainer()-get(Producer::class); $result $producer-produce($message);上述示例直接使用ApplicationContext获取Hyperf\Amqp\Producer并不符合规范DI 容器的具体用法请参见依赖注入章节。投递成功后produce()方法默认返回true。从源码看Producer::produce(ProducerMessageInterface $producerMessage, bool $confirm false, int $timeout 5)支持两个可选参数$confirm开启发布确认模式Publisher Confirms此时会等待 Broker 的 ack 确认后才返回true$timeout为等待确认的超时秒数。消息序列化由 Packer 完成生产者的默认properties包含content_type text/plain与持久化投递模式DELIVERY_MODE_PERSISTENT见 ProducerMessage.php。消费消息Consumer使用gen:amqp-consumer命令创建一个消费者php bin/hyperf.php gen:amqp-consumer DemoConsumer在DemoConsumer中可以修改#[Consumer]注解对应的字段来替换exchange、routingKey和queue。其中$data就是反序列化后的消息数据。使用#[Consumer]注解时需引入use Hyperf\Amqp\Annotation\Consumer;命名空间。?php declare(strict_types1); namespace App\Amqp\Consumers; use Hyperf\Amqp\Annotation\Consumer; use Hyperf\Amqp\Message\ConsumerMessage; use Hyperf\Amqp\Result; use PhpAmqpLib\Message\AMQPMessage; #[Consumer(exchange: hyperf, routingKey: hyperf, queue: hyperf, nums: 1)] class DemoConsumer extends ConsumerMessage { public function consumeMessage($data, AMQPMessage $message): Result { print_r($data); return Result::ACK; } }从 Annotation/Consumer.php 源码可见#[Consumer]注解支持exchange、routingKey支持数组、queue、name进程名称默认Consumer、nums消费进程数量、enable是否随服务启动、maxConsumption最大消费条数与pool共 8 个参数。消费者注册由 ConsumerManager.php 配合 MainWorkerStartListener.php 完成默认情况下使用了#[Consumer]注解后框架会为每个消费者自动创建独立子进程并常驻消费当子进程异常退出时还会自动重新拉起。禁止消费进程自启开发阶段调试消费者时若不想让消费者随服务启动而自动消费其他消息可以通过两种方式禁用在#[Consumer]注解中配置enablefalse默认为true即跟随服务启动在消费者类中重写isEnable()方法返回false。?php declare(strict_types1); namespace App\Amqp\Consumers; use Hyperf\Amqp\Annotation\Consumer; use Hyperf\Amqp\Message\ConsumerMessage; use Hyperf\Amqp\Result; use PhpAmqpLib\Message\AMQPMessage; #[Consumer(exchange: hyperf, routingKey: hyperf, queue: hyperf, nums: 1, enable: false)] class DemoConsumer extends ConsumerMessage { public function consumeMessage($data, AMQPMessage $message): Result { print_r($data); return Result::ACK; } public function isEnable(): bool { return parent::isEnable(); } }isEnable()方法定义在 ConsumerMessage.php默认返回$this-enable重写时可直接返回false或基于业务条件如当前环境动态判断。设置最大消费数修改#[Consumer]注解中的maxConsumption属性可设置该消费者最大处理的消息数。从 Consumer.php 源码看消费循环内部会维护一个$currentConsumption计数器当达到maxConsumption时主动跳出消费循环并关闭连接随后由进程管理机制重启消费者进程实现消费 N 条后自动重启的内存泄漏兜底策略。设置并发消费影响消费速率的参数有三处#[Consumer]注解的nums开启多个消费者进程每个进程独立消费ConsumerMessage基类下的$qos属性通过重写其中的prefetch_size或prefetch_count控制每次从服务端预取的消息数量配置文件中的concurrent.limit控制单个消费者进程内并发消费协程的最大数量见上文getConcurrent()源码解析。消费结果框架会根据consumeMessage()方法返回的Result来决定对消息的响应行为。Result枚举定义在 src/amqp/src/Result.php共 4 种取值返回值行为\Hyperf\Amqp\Result::ACK确认消息被正确消费通过basic_ack响应\Hyperf\Amqp\Result::NACK消息未被正确消费通过basic_nack响应\Hyperf\Amqp\Result::REQUEUE消息未被正确消费通过basic_reject响应并使消息重新入列\Hyperf\Amqp\Result::DROP消息未被正确消费通过basic_reject响应丢弃对应的底层行为可在 Consumer.php 的getCallback()回调中确认ACK调用basic_ack($deliveryTag)NACK调用basic_nack($deliveryTag, false, $consumerMessage-isRequeue())REQUEUE在isRequeue()为真时调用basic_reject($deliveryTag, true)让消息重回队列其余情况调用basic_reject($deliveryTag, false)直接拒绝。另外若consumeMessage()抛出异常框架会分发FailToConsume事件并记录错误日志最终按Result::DROP处理。QOS 配置消费者基类默认的$qos为[prefetch_size 0, prefetch_count 1, global false]见 ConsumerMessage.php可以按需重写?php declare(strict_types1); namespace App\Amqp\Consumers; use Hyperf\Amqp\Annotation\Consumer; use Hyperf\Amqp\Message\ConsumerMessage; use Hyperf\Amqp\Result; use PhpAmqpLib\Message\AMQPMessage; #[Consumer(exchange: hyperf, routingKey: hyperf, queue: hyperf, nums: 1)] class DemoConsumer extends ConsumerMessage { protected ?array $qos [ // AMQP 默认并没有实现此配置 prefetch_size 0, // 同一个消费者最高同时可以处理的消息数 prefetch_count 30, // 因为 Hyperf 默认一个 Channel 只消费一个队列所以 global 设置为 true/false 效果是一样的 global false, ]; public function consumeMessage($data, AMQPMessage $message): Result { print_r($data); return Result::ACK; } }在 Consumer.php 的declare()方法中$qos数组中的prefetch_size、prefetch_count、global会被提取并调用basic_qos()应用到 Channel从而控制服务端推送消息的窗口大小。根据环境自定义消费进程数量#[Consumer]注解中的nums属性用于设置消费进程数量若需根据环境动态调整可以重写getNums()方法#[Consumer( exchange: hyperf, routingKey: hyperf, queue: hyperf, name: hyperf, nums: 1 )] final class DemoConsumer extends ConsumerMessage { public function getNums(): int { if (is_debug()) { return 10; } return parent::getNums(); } }getNums()基类实现见 ConsumerMessage.php返回$this-nums此处通过is_debug()判断调试环境后返回更多进程实现开发/生产环境消费能力的差异化。延迟队列AMQP 的延迟队列基于rabbitmq-delayed-message-exchange插件实现需要注意该延迟队列并不会根据延迟时间进行排序。一旦先投递了一个延迟 10s 的任务再向同一队列投递一个延迟 5s 的任务那么第二个 5s 任务也一定会在第一个 10s 任务完成后才会被消费。因此需要根据时间粒度拆分不同的队列如果想要更灵活支持按时间排序的延迟队列可以尝试将异步队列与 AMQP 配合使用。使用延迟队列前需要先下载并激活 RabbitMQ 延迟消息插件wget https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases/download/3.9.0/rabbitmq_delayed_message_exchange-3.9.0.ez cp rabbitmq_delayed_message_exchange-3.9.0.ez /opt/rabbitmq/plugins/ rabbitmq-plugins enable rabbitmq_delayed_message_exchange延迟生产者使用命令创建生产者这里以direct类型举例fanout、topic类型只需修改生产者和消费者中的type即可php bin/hyperf.php gen:amqp-producer DelayDirectProducer在DelayDirectProducer中引入ProducerDelayedMessageTrait?php namespace App\Amqp\Producer; use Hyperf\Amqp\Annotation\Producer; use Hyperf\Amqp\Message\ProducerDelayedMessageTrait; use Hyperf\Amqp\Message\ProducerMessage; use Hyperf\Amqp\Message\Type; #[Producer] class DelayDirectProducer extends ProducerMessage { use ProducerDelayedMessageTrait; protected string $exchange ext.hyperf.delay; protected Type|string $type Type::DIRECT; protected array|string $routingKey ; public function __construct($data) { $this-payload $data; } }从 ProducerDelayedMessageTrait.php 源码可以看到该 Trait 的核心实现setDelayMs(int $millisecond, string $name x-delay)将延迟毫秒数写入消息的application_headersAMQPTable供插件识别延迟时间重写getExchangeBuilder()将交换机类型设置为x-delayed-message并通过参数x-delayed-type声明底层实际的交换机类型如direct/fanout/topic实现延迟 任意类型的组合。延迟消费者php bin/hyperf.php gen:amqp-consumer DelayDirectConsumer在DelayDirectConsumer中同时引入ProducerDelayedMessageTrait与ConsumerDelayedMessageTrait?php declare(strict_types1); namespace App\Amqp\Consumer; use Hyperf\Amqp\Annotation\Consumer; use Hyperf\Amqp\Message\ConsumerDelayedMessageTrait; use Hyperf\Amqp\Message\ConsumerMessage; use Hyperf\Amqp\Message\ProducerDelayedMessageTrait; use Hyperf\Amqp\Message\Type; use Hyperf\Amqp\Result; use PhpAmqpLib\Message\AMQPMessage; #[Consumer(nums: 1)] class DelayDirectConsumer extends ConsumerMessage { use ProducerDelayedMessageTrait; use ConsumerDelayedMessageTrait; protected string $exchange ext.hyperf.delay; protected string $queue queue.hyperf.delay; protected Type|string $type Type::DIRECT; //Type::FANOUT; protected array|string $routingKey ; public function consumeMessage($data, AMQPMessage $message): Result { var_dump($data, delaydirect consumeTime: . (microtime(true))); return Result::ACK; } }注意延迟消费者中也引入了ProducerDelayedMessageTrait其目的是让消费者声明出与生产者一致的x-delayed-message交换机确保队列绑定成功。ConsumerDelayedMessageTrait.php 则重写了getQueueBuilder()为队列声明死信交换机参数x-dead-letter-exchange delayed——延迟到期的消息会通过死信机制被路由到真实消费队列。生产延迟消息以下代码在 Command 中演示如何发送延迟消息具体用法以实际业务为准?php declare(strict_types1); namespace App\Command; use App\Amqp\Producer\DelayDirectProducer; //use App\Amqp\Producer\DelayFanoutProducer; //use App\Amqp\Producer\DelayTopicProducer; use Hyperf\Amqp\Producer; use Hyperf\Command\Annotation\Command; use Hyperf\Command\Command as HyperfCommand; use Hyperf\Context\ApplicationContext; use Psr\Container\ContainerInterface; #[Command] class DelayCommand extends HyperfCommand { protected ContainerInterface $container; public function __construct(ContainerInterface $container) { $this-container $container; parent::__construct(demo:command); } public function configure() { parent::configure(); $this-setDescription(Hyperf Demo Command); } public function handle() { // 1. delayed direct $message new DelayDirectProducer(delaydirect produceTime:.(microtime(true))); // 2. delayed fanout //$message new DelayFanoutProducer(delayfanout produceTime:.(microtime(true))); // 3. delayed topic //$message new DelayTopicProducer(delaytopic produceTime: . (microtime(true))); $message-setDelayMs(5000); $producer ApplicationContext::getContainer()-get(Producer::class); $producer-produce($message); } }执行命令行即可生产一条延迟 5 秒的消息php bin/hyperf.php demo:commandRPC 远程过程调用除了典型的消息队列场景还可以通过 AMQP 实现 RPC 远程过程调用组件为此提供了完整支持生产者发送请求消息携带reply_to与correlation_id消费者处理后将结果通过reply()方法写回临时响应队列生产者同步等待结果。创建 RPC 消费者RPC 消费者与普通消费者的实现基本一致唯一区别是需要通过reply()方法将数据返回给生产者?php declare(strict_types1); namespace App\Amqp\Consumer; use Hyperf\Amqp\Annotation\Consumer; use Hyperf\Amqp\Message\ConsumerMessage; use Hyperf\Amqp\Result; use PhpAmqpLib\Message\AMQPMessage; #[Consumer(exchange: hyperf, routingKey: hyperf, queue: rpc.reply, name: ReplyConsumer, nums: 1, enable: true)] class ReplyConsumer extends ConsumerMessage { public function consumeMessage($data, AMQPMessage $message): Result { $data[message] . Reply: . $data[message]; $this-reply($data, $message); return Result::ACK; } }reply()方法实现在 ConsumerMessage.php它取出请求消息中的reply_to队列与correlation_id将处理结果打包成新的AMQPMessage发布到默认交换机空字符串上的reply_to队列从而把结果路由回发起 RPC 的生产者。发起 RPC 调用作为调用方发起一次 RPC 调用非常简单只需从 DI 容器获得Hyperf\Amqp\RpcClient并调用call()方法返回值即为消费者reply回来的数据?php use Hyperf\Amqp\Message\DynamicRpcMessage; use Hyperf\Amqp\RpcClient; use Hyperf\Context\ApplicationContext; $rpcClient ApplicationContext::getContainer()-get(RpcClient::class); // 在 DynamicRpcMessage 上设置与 Consumer 一致的 Exchange 和 RoutingKey $result $rpcClient-call(new DynamicRpcMessage(hyperf, hyperf, [message Hello Hyperf])); // $result: // array(1) { // [message] // string(18) Reply:Hello Hyperf // }从 RpcClient.php 源码看call()的执行流程为按pool exchange queue维度从 Channel 池取出空闲 RpcChannel池上限默认 64可通过构造参数$maxChannels调整→ 声明临时响应队列并basic_consume监听 → 携带correlation_id与reply_to发布请求消息 → 通过协程 Channel 等待响应默认超时 5 秒超时抛出TimeoutException→ 反序列化响应体返回并将 Channel 归还池中复用。同时RpcChannel 在收到响应时会校验correlation_id是否匹配避免并发调用时结果串线。抽象 RpcMessage上面的调用过程直接通过DynamicRpcMessage定义 Exchange 和 RoutingKey 并传递数据。在生产项目中更推荐对RpcMessage做一层抽象统一 Exchange 和 RoutingKey 的定义避免每次调用都重复声明?php use Hyperf\Amqp\Message\RpcMessage; class FooRpcMessage extends RpcMessage { protected string $exchange hyperf; protected array|string $routingKey hyperf; public function __construct($data) { // 要传递的数据 $this-payload $data; } }此后发起 RPC 调用时只需直接传入FooRpcMessage实例$result $rpcClient-call(new FooRpcMessage([message Hello Hyperf]));无需每次调用都重复定义 Exchange 和 RoutingKey消息结构也得以统一管理便于在多个服务间沉淀为公共契约。小结围绕hyperf/amqp组件本文完整覆盖了从安装配置、Producer 投递、Consumer 消费含进程自启控制、最大消费数、并发调优、消费结果语义、QOS 预取到延迟队列与消息 RPC 的全链路用法。结合 src/amqp/src 下的源码可以看到组件的核心类如 Producer.php、Consumer.php、RpcClient.php 与 ConsumerMessage.php 均围绕注解驱动 协程并发 连接池复用的设计展开实际项目中可直接参考 src/amqp/publish/amqp.php 的默认配置进行环境化改造也可以阅读 src/amqp/tests 下的测试用例加深对组件行为的理解。赞分享后端微服务【免费下载链接】hyperf A coroutine framework that focuses on hyperspeed and flexibility. Building microservice or middleware with ease.项目地址https://gitcode.com/gh_mirrors/hy/hyperf点击查看免费下载相关推荐Hyperf AMQP 组件实战指南基于 RabbitMQ 的生产消费、延迟队列与 RPC 调用Hyperf AMQP 组件实战指南基于 RabbitMQ 的生产消费、延迟队列与 RPC 调用 Hyperf 的 hyperf/amqp https://l后端Web框架微服务RPC框架异步编程Hyperf AMQP 组件实战指南RabbitMQ 消息生产、消费、延迟队列与 RPC 调用Hyperf AMQP 组件实战指南RabbitMQ 消息生产、消费、延迟队列与 RPC 调用 本指南以 Hyperf 官方文档《Komponen AMQP》后端微服务Hyperf AMQP 组件实战指南生产者、消费者、延迟队列与基于 AMQP 的 RPCHyperf AMQP 组件实战指南生产者、消费者、延迟队列与基于 AMQP 的 RPC 本文以 Hyperf 框架内置的 hyperf/amqp 组件基于后端Web框架微服务RPC框架异步编程上一篇Swagger Petstore Bash 客户端 UserApi 使用指南基于 swagger-codegen 生成的用户管理接口实战下一篇Trivy 自定义检查的 Input Schema 详解为 Rego 检查启用输入类型校验创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑