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

使用php-amqplib连接ActiveMQ Classic Docker镜像遇AMQP版本兼容错误

问题:使用php-amqplib连接ActiveMQ Classic消费消息时出现AMQP版本不兼容错误

错误现象

执行sail artisan horizon:consume-activemq后抛出异常:

PhpAmqpLib\Exception\AMQPInvalidFrameException

Invalid frame type 65

...

7 app/Queue/Connectors/ActiveMQConnector.php:15
PhpAmqpLib\Connection\AMQPStreamConnection::__construct()

8 app/Console/Commands/ConsumeActiveMQMessages.php:29
App\Queue\Connectors\ActiveMQConnector::connect()

查看Docker日志,发现核心原因:

Connection attempt from non AMQP v1.0 client. AMQP,0,0,9,1
2024-02-22 13:24:44 WARN | Transport Connection to: tcp://192.168.65.1:32999 failed: org.apache.activemq.transport.amqp.AmqpProtocolException: Connection from the client using unsupported AMQP attempted

核心问题

php-amqplib是针对AMQP 0-9-1协议开发的客户端库,而当前连接的ActiveMQ Classic的AMQP端口(5672)默认要求使用AMQP 1.0协议,两者协议版本不兼容,导致连接失败。

配置信息

queue.php

'connections' => [
    ...
    'activemq' => [
        'driver' => 'activemq',
        'host' => env('ACTIVEMQ_HOST', 'localhost'),
        'port' => env('ACTIVEMQ_PORT', 61613),
        'username' => env('ACTIVEMQ_USERNAME', 'guest'),
        'password' => env('ACTIVEMQ_PASSWORD', 'guest'),
        'queue' => env('ACTIVEMQ_QUEUE', ''),
        'exchange_name' => env('ACTIVEMQ_EXCHANGE_NAME', ''),
    ],
]

.env文件

ACTIVEMQ_HOST=host.docker.internal
ACTIVEMQ_PORT=5672
ACTIVEMQ_USER=admin
ACTIVEMQ_PASSWORD=admin
ACTIVEMQ_QUEUE=activemqTest

ActiveMQServiceProvider.php

<?php

namespace App\Providers;

use App\Queue\Connectors\ActiveMQConnector;
use Illuminate\Queue\QueueManager;
use Illuminate\Support\ServiceProvider;

class ActiveMQServiceProvider extends ServiceProvider
{
    /**
     * Register services.
     */
    public function register(): void
    {
    }

    /**
     * Bootstrap services.
     */
    public function boot(): void
    {
        $this->app->make(QueueManager::class)->addConnector('activemq', function () {
            return new ActiveMQConnector();
        });
    }
}

ActiveMQConnector.php

<?php

namespace App\Queue\Connectors;

use Illuminate\Queue\Connectors\ConnectorInterface;
use PhpAmqpLib\Connection\AMQPStreamConnection;

class ActiveMQConnector implements ConnectorInterface
{
    /**
     * @throws \Exception
     */
    public function connect(array $config)
    {
        return new AMQPStreamConnection(
            $config['host'],
            $config['port'],
            $config['username'],
            $config['password'],
            $config['vhost']
        );
    }
}

ConsumeActiveMQMessages.php

<?php

namespace App\Console\Commands;

use App\Queue\Connectors\ActiveMQConnector;
use Illuminate\Console\Command;
use PhpAmqpLib\Message\AMQPMessage;

class ConsumeActiveMQMessages extends Command
{
    protected $signature = 'horizon:consume-activemq';

    protected $description = 'Consume messages from ActiveMQ and process them within Horizon';

    /**
     * @throws \Exception
     */
    public function handle()
    {
        $connector = new ActiveMQConnector();
        $config = [
            'host' => config('queue.connections.activemq.host'),
            'port' => config('queue.connections.activemq.port'),
            'username' => config('queue.connections.activemq.username'),
            'password' =>config('queue.connections.activemq.password'),
            'vhost' => config('queue.connections.activemq.vhost') !== null ?config('queue.connections.activemq.vhost') : '/',
        ];

        $connection = $connector->connect($config);

        $channel = $connection->channel();

        $callback = function (AMQPMessage $message) {
            $this->processMessage($message);
        };

        $channel->basic_consume(config('queue.connections.activemq.queue'), '', false, true, false, false, $callback);

        while ($channel->is_consuming()) {
            $channel->wait();
        }

        $channel->close();
        $connection->close();
    }

    protected function processMessage(AMQPMessage $message)
    {
        $this->info('Received message: ' . $message->getBody());
    }
}

环境信息

  • Laravel V10.10
  • PHP V8.1

解决方案

有两种可行方案:

方案1:修改ActiveMQ配置,启用AMQP 0-9-1兼容模式

在ActiveMQ的activemq.xml配置文件中,修改AMQP传输适配器,添加amqp091Compatible="true"属性,让ActiveMQ的5672端口兼容AMQP 0-9-1协议:

<transportConnector name="amqp" uri="amqp://0.0.0.0:5672?maximumConnections=1000&amp;wireFormat.maxFrameSize=104857600&amp;amqp091Compatible=true"/>

修改后重启ActiveMQ服务,保持现有代码和配置不变即可连接。

方案2:改用支持AMQP 1.0的PHP客户端库

放弃php-amqplib,使用支持AMQP 1.0的库,比如php-amqp/php-amqp1.0。步骤如下:

  1. 安装依赖:
composer require php-amqp/php-amqp1.0
  1. 修改ActiveMQConnector.php,适配AMQP 1.0连接:
<?php

namespace App\Queue\Connectors;

use Illuminate\Queue\Connectors\ConnectorInterface;
use PhpAmqp\Client\Connection;
use PhpAmqp\Client\ConnectionOptions;

class ActiveMQConnector implements ConnectorInterface
{
    public function connect(array $config)
    {
        $options = new ConnectionOptions();
        $options->setHost($config['host'])
                ->setPort($config['port'])
                ->setUsername($config['username'])
                ->setPassword($config['password']);

        return new Connection($options);
    }
}
  1. 更新消费命令的逻辑,适配AMQP 1.0的消息消费API。

额外注意点

  • 若之前尝试切换到61613端口,该端口是ActiveMQ的STOMP协议端口,php-amqplib不支持,因此会失败。
  • 方案1适合快速兼容现有代码,方案2是长期的标准AMQP 1.0适配方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 15:27:02