使用php-amqplib连接ActiveMQ Classic Docker镜像遇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&wireFormat.maxFrameSize=104857600&amqp091Compatible=true"/>
修改后重启ActiveMQ服务,保持现有代码和配置不变即可连接。
方案2:改用支持AMQP 1.0的PHP客户端库
放弃php-amqplib,使用支持AMQP 1.0的库,比如php-amqp/php-amqp1.0。步骤如下:
- 安装依赖:
composer require php-amqp/php-amqp1.0
- 修改
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); } }
- 更新消费命令的逻辑,适配AMQP 1.0的消息消费API。
额外注意点
- 若之前尝试切换到61613端口,该端口是ActiveMQ的STOMP协议端口,php-amqplib不支持,因此会失败。
- 方案1适合快速兼容现有代码,方案2是长期的标准AMQP 1.0适配方案。
内容的提问来源于stack exchange,提问作者Ali Ibrahim

