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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 00:11:04