You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Laravel切换laravel-amqp后,如何持续消费RabbitMQ消息?

解决方案

核心结论

ssi-anik/laravel-amqp 包没有内置等效于rabbitmq:consume或queue:work的Artisan消费命令,你需要自定义Artisan命令来实现持续监听并消费队列消息的功能。

自定义消费命令步骤

  1. 生成自定义命令
    在终端执行以下命令生成基础命令文件:

    php artisan make:command AmqpConsume
    
  2. 编写命令逻辑
    打开生成的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,
            ]);
        }
    }
    
  3. 注册命令(可选)
    Laravel 9会自动发现Console目录下的命令,若未自动加载,可手动在app/Console/Kernel.php的$commands数组中添加:

    protected $commands = [
        \App\Console\Commands\AmqpConsume::class,
    ];
    
  4. 运行消费命令
    在终端执行自定义命令,指定要消费的队列名称:

    php artisan amqp:consume your-queue-name
    

关键说明

  • 代码采用手动消息确认机制,确保只有消息处理成功后才会从队列移除,避免消息丢失。
  • 加入信号处理器,支持通过Ctrl+C或系统信号优雅停止消费进程。
  • 可根据业务需求调整消息处理逻辑、异常处理策略(如重试机制、死信队列配置等)。

内容的提问来源于stack exchange,提问作者Sam

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 09:52:35