PHP版AMQP连接正常,Node.js版配置适配问题求助
问题
现有一段可正常连接外部AMQP服务器的PHP代码:
<?php require_once '/usr/share/php/PhpAmqpLib/autoload.php'; use PhpAmqpLib\Connection\AMQPStreamConnection; $connection = new AMQPStreamConnection('hostname', 5672, 'username', 'password', "vhost", false, 'AMQPLAIN', null, 'en_US', 1160, 1160, null, false, 580); $channel = $connection->channel(); $queue = 'queue'; $channel->basic_qos(0,1000,false); $callback = function($msg) { #file_put_contents('test.log', "\n===> NEW Message start \n", FILE_APPEND); #file_put_contents('test.log', print_r($msg->body, true), FILE_APPEND); #file_put_contents('test.log', "\n===> NEW Message end \n", FILE_APPEND); echo "=====> Received start\n"; echo " [x] Received ", $msg->body, "\n"; echo "=====> Received end\n"; }; $channel->basic_consume($queue, 'consumer', false, true, false, false, $callback); while(count($channel->callbacks)) { $channel->wait(); } $channel->close(); $connection->close() ?>
需要将该功能迁移至Node.js环境,使用amqplib库,但当前配置无法适配PHP中的连接参数,现有Node.js代码如下:
const amqp = require('amqplib/callback_api'); const config = { protocol: 'amqp', hostname: 'hostname', port: 5672, username: 'username', password: 'password', locale: 'en_US', vhost: 'vhost', } amqp.connect(config, function (error0, connection) { if (error0) { throw error0; } connection.createChannel(function (error1, channel) { if (error1) { throw error1; } var queue = 'queue'; channel.assertQueue(queue, { durable: false }); console.log(" [*] Waiting for messages in %s. To exit press CTRL+C", queue); channel.consume(queue, function (msg) { console.log(" [x] Received %s", msg.content.toString()); }, { noAck: true }); }); });
请调整Node.js代码配置,使其实现与PHP代码完全一致的AMQP连接和消费功能。
调整后的Node.js代码
const amqp = require('amqplib/callback_api'); const config = { protocol: 'amqp', hostname: 'hostname', port: 5672, username: 'username', password: 'password', locale: 'en_US', vhost: 'vhost', // 对应PHP的连接超时参数(580秒转毫秒) timeout: 580000, // 与PHP一致的心跳间隔 heartbeat: 1160, // 匹配PHP指定的AMQPLAIN认证机制 authMechanism: 'AMQPLAIN' }; amqp.connect(config, function (error0, connection) { if (error0) { throw error0; } connection.createChannel(function (error1, channel) { if (error1) { throw error1; } const queue = 'queue'; const consumerTag = 'consumer'; // 对应PHP的consumer标识 // 复制PHP的QoS配置:一次最多接收1000条未确认消息 channel.prefetch(1000, false); // 断言队列,与PHP隐式行为保持一致 channel.assertQueue(queue, { durable: false }); console.log("=====> Waiting for messages in %s. To exit press CTRL+C", queue); // 消费消息,同步PHP的参数与输出格式 channel.consume(queue, function (msg) { console.log("=====> Received start"); console.log(" [x] Received %s", msg.content.toString()); console.log("=====> Received end"); }, { noAck: true, // 对应PHP的auto_ack=true consumerTag: consumerTag // 指定消费者标签 }); }); });
关键配置匹配说明
- 超时与心跳:PHP中
580是连接超时(秒),Node.js需转为毫秒单位的timeout;心跳参数1160直接对应Node.js的heartbeat配置。 - 认证机制:PHP指定
AMQPLAIN认证,Node.js通过authMechanism参数同步该配置。 - 消费者标签:PHP的
basic_consume第二个参数是消费者标识consumer,Node.js在consume的选项中通过consumerTag指定。 - QoS规则:PHP的
basic_qos(0,1000,false)对应Node.js的channel.prefetch(1000, false),控制消息预取数量。 - 输出格式:调整console输出内容,与PHP的日志格式完全一致。
内容的提问来源于stack exchange,提问作者tai kim
相关产品推荐
相关产品推荐

