Laravel切换laravel-amqp后,如何持续消费RabbitMQ消息?
解决方案
核心结论
ssi-anik/laravel-amqp 包没有内置等效于rabbitmq:consume或queue:work的Artisan消费命令,你需要自定义Artisan命令来实现持续监听并消费队列消息的功能。
自定义消费命令步骤
生成自定义命令
在终端执行以下命令生成基础命令文件:php artisan make:command AmqpConsume编写命令逻辑
打开生成的app/Console/Commands/AmqpConsume.php文件,替换内容为以下适配Laravel 9及对应包版本的代码:<?php namespace App\Console\Commands; use Anik\Amqp\Consumers\Consumer; use Anik\Amqp\Facades\Amqp; use Illuminate\Console\Command; use PhpAmqpLib\Message\AMQPMessage; class AmqpConsume extends Command { /** * 命令的名称和签名 * * @var string */ protected $signature = 'amqp:consume {queue : 要消费的队列名称}'; /** * 命令描述 * * @var string */ protected $description = '持续监听并消费指定的AMQP队列消息'; /** * 执行命令 */ public function handle() { $queueName = $this->argument('queue'); // 注册信号处理器,实现优雅停止 pcntl_signal(SIGINT, function () { $this->info('正在停止消费...'); exit(0); }); pcntl_signal(SIGTERM, function () { $this->info('正在停止消费...'); exit(0); }); $this->info("开始监听队列: {$queueName}"); // 使用包的消费者持续消费 Amqp::connection('queuename')->consume($queueName, function (AMQPMessage $message, Consumer $consumer) { try { // 处理消息逻辑,根据业务需求修改 $payload = json_decode($message->body, true); $this->info("收到消息: " . json_encode($payload)); // 确认消息已处理完成 $message->ack(); } catch (\Exception $e) { $this->error("处理消息失败: " . $e->getMessage()); // 处理失败时拒绝消息,可根据业务选择是否重新入队 $message->nack(false, false); } // 处理信号,避免进程阻塞无法停止 pcntl_signal_dispatch(); }, [ 'consumer_tag' => 'laravel-amqp-consumer', 'no_local' => false, 'no_ack' => false, // 手动确认消息,需调用$message->ack() 'exclusive' => false, ]); } }注册命令(可选)
Laravel 9会自动发现Console目录下的命令,若未自动加载,可手动在app/Console/Kernel.php的$commands数组中添加:protected $commands = [ \App\Console\Commands\AmqpConsume::class, ];运行消费命令
在终端执行自定义命令,指定要消费的队列名称:php artisan amqp:consume your-queue-name
关键说明
- 代码采用手动消息确认机制,确保只有消息处理成功后才会从队列移除,避免消息丢失。
- 加入信号处理器,支持通过
Ctrl+C或系统信号优雅停止消费进程。 - 可根据业务需求调整消息处理逻辑、异常处理策略(如重试机制、死信队列配置等)。
内容的提问来源于stack exchange,提问作者Sam
相关产品推荐
相关产品推荐

