【问题标题】:How to Consume & Publish message from Amazon MQ in Nodejs?如何在 Node Js 中使用和发布来自 Amazon MQ 的消息?
【发布时间】:2018-07-08 09:28:09
【问题描述】:

我需要使用 Amazon MQ 在 Nodejs 中使用 amqp 协议来消费和发布消息到队列。我已经设置了 AWS MQ,定义了代理并创建了一个队列。

我已经关注了 AWS Javascript SDK,但我仍然找不到任何方法来消费和发布消息到队列。

谁能帮助我如何使用 amqp 协议连接到 AWS MQ 以消费和发布消息到队列。

谢谢

【问题讨论】:

    标签: node.js amqp amazon-mq


    【解决方案1】:

    我使用了 amqp10 npm 模块,该模块用于消费并将消息发布到 AWS MQ。

    下面是代码:

      const AMQPClient = require('amqp10').Client;
      const Policy = require('amqp10').Policy;
    

    1.从AWS MQ消费消息:

          let client = new AMQPClient(Policy.Utils.RenewOnSettle(1, 1, 
    Policy.ServiceBusQueue));
          let connectionString = 'your_connnection_string';
          client.connect(connectionString)
              .then(function() {
                  console.log("Connected");
                  return Promise.all([
    
    client.createReceiver(configurationHolder.config.getMessageQueueName)
                  ]);
              })
              .spread(function(receiver) {
                      receiver.on('errorReceived', function(rx_err) {
                          console.warn('===> RX ERROR: ', rx_err);
                          return err;
                      });
                      receiver.on('message', function(message) {
                          client.disconnect().then(function() {
                          console.log('disconnected, when we get the message from the queue);
                          return message.body;
                      });
                  });
              })
              .error(function(e) {
                      console.warn('connection error: ', e);
                      return err;
                  });
    
    1. 向 AWS MQ 发布消息:

      let client = new AMQPClient(Policy.merge({
          senderLinkPolicy: {
              callbackPolicy: Policy.Utils.SenderCallbackPolicies.OnSent
          }
      }, Policy.DefaultPolicy));
      
      
          client.connect(connectionString, {
                  'saslMechanism': 'ANONYMOUS'
              })
              .then(function() {
                  console.log("Connected");
                  return Promise.all([
                      client.createSender(queueName)
                  ]);
              })
              .spread(function(sender) {
                  sender.on('errorReceived', function(tx_err) {
                      console.warn('===> TX ERROR: ', tx_err);
                      return err;
                  });
                  var options = {
                      annotations: {
                          'x-opt-partition-key': 'pk' + msgValue
                      }
                  };
                  return sender.send(JSON.stringify(msgValue), 
      options).then(function(state) {
                      client.disconnect().then(function() {
                          console.log('disconnected, when we saw the value we 
      inserted after publish to AWS MQ.');
                          return state;
                      });
                  });
              })
              .error(function(e) {
                  console.warn('connection error: ', e);
                  return err;
              });
      

    谢谢

    【讨论】:

    • 嘿 Ankit 你试过这种高负载的方法(库)吗?我还看到您首先连接然后创建发件人,但是当您发送数百万条消息时,您是否始终创建连接然后发件人?
    • 嗨 Shark,我只在第一条消息发布时创建连接和发送方(特定于队列)。虽然我没有找到任何发布批量消息的规定,因为它会逐条发布每条消息,而且肯定更快,因为我不需要为每条消息创建连接和发件人。
    猜你喜欢
    • 2022-07-06
    • 1970-01-01
    • 2021-07-23
    • 2019-05-30
    • 2016-12-02
    • 2019-08-25
    • 1970-01-01
    • 2014-06-21
    • 1970-01-01
    相关资源
    最近更新 更多