【问题标题】:Nodejs sqs queue processorNodejs sqs队列处理器
【发布时间】:2013-07-16 23:40:53
【问题描述】:

我正在尝试编写一个 nodejs sqs 队列处理器。

"use strict";
var appConf = require('./config/appConf');
var AWS = require('aws-sdk');
AWS.config.loadFromPath('./config/aws_config.json');
var sqs = new AWS.SQS();
var exec = require('child_process').exec;
function readMessage() {
  sqs.receiveMessage({
    "QueueUrl": appConf.sqs_distribution_url,
    "MaxNumberOfMessages": 1,
    "VisibilityTimeout": 30,
    "WaitTimeSeconds": 20
  }, function (err, data) {
    var sqs_message_body;
    if (data.Messages) {
      if (typeof data.Messages[0] !== 'undefined' && typeof data.Messages[0].Body !== 'undefined') {
        //sqs msg body
        sqs_message_body = JSON.parse(data.Messages[0].Body);
        //make call to nodejs handler in codeigniter
        exec('php '+ appConf.CI_FC_PATH +'/index.php nodejs_handler make_contentq_call "'+ sqs_message_body.contentq_cat_id+'" "'+sqs_message_body.cnhq_cat_id+'" "'+sqs_message_body.network_id+'"',
          function (error, stdout, stderr) {
            if (error) {
              throw error;
            }
            console.log('stdout: ' + stdout);
            if(stdout == 'Success'){
              //delete message from queue
              sqs.deleteMessage({
                "QueueUrl" : appConf.sqs_distribution_url,
                "ReceiptHandle" :data.Messages[0].ReceiptHandle
              });
            }
          });
      }
    }
  });
}
readMessage();

以上代码适用于队列中的单个消息。我应该如何编写此脚本,以便在处理完所有消息之前一直轮询队列中的消息?我应该使用设置超时吗?

【问题讨论】:

    标签: node.js amazon-sqs


    【解决方案1】:

    如果您使用节点,请使用https://www.npmjs.com/package/sqs-worker 模块。它会为你完成这项工作。

    var SQSWorker = require('sqs-worker')
    
    var options =
     { url: 'https://sqs.eu-west-1.amazonaws.com/001123456789/my-queue'
    }
    
    var queue = new SQSWorker(options, worker)
    
    function worker(notifi, done) {
      var message;
      try {
        message = JSON.parse(notifi.Data)
      } catch (err) {
        throw err
      }
    
       // Do something with `message` 
    
       var success = true
    
       // Call `done` when you are done processing a message. 
       // If everything went successfully and you don't want to see it any more, 
       // set the second parameter to `true`. 
       done(null, success)
    }
    

    【讨论】:

    【解决方案2】:

    首先,您应该明确地使用亚马逊提供的 long polling 技术,据我了解,您已经在使用它,因为您在 sqs.receiveMessage 调用中有 "WaitTimeSeconds": 20 参数。希望大家不要忘记在AWS Web interface中进行配置。

    关于消息轮询 - 您可以使用不同的技术,包括计时器,但我认为最简单的方法是在 receiveMessage(甚至是 exec)回调的末尾调用您的 readMessage() 函数功能。因此队列中下一条消息的处理(或等待)将在队列中上一条消息的处理结束后立即开始。

    更新:

    至于我,在您的新代码版本中有很多 readMessage() 调用。我认为最好将其最小化以使代码更清晰且易于维护。但是,例如,如果您在主 receiveMessage 回调结束时留下唯一的一个调用,您将收到许多并行运行的 PHP 工作脚本 - 从性能的角度来看,这可能还不错 -但是您必须添加一些复杂的脚本来控制并行工作人员的数量。我认为您可以在exec 回调中减少一些呼叫,尝试加入ifs 并在主回调中加入呼叫。

    "use strict";
    var appConf = require('./config/appConf');
    var AWS = require('aws-sdk');
    AWS.config.loadFromPath('./config/aws_config.json');
    var delay = 20 * 1000;
    var sqs = new AWS.SQS();
    var exec = require('child_process').exec;
    function readMessage() {
      sqs.receiveMessage({
        "QueueUrl": appConf.sqs_distribution_url,
        "MaxNumberOfMessages": 1,
        "VisibilityTimeout": 30,
        "WaitTimeSeconds": 20
      }, function (err, data) {
        var sqs_message_body;
        if (data.Messages) 
          && (typeof data.Messages[0] !== 'undefined' && typeof data.Messages[0].Body !== 'undefined')) {
            //sqs msg body
            sqs_message_body = JSON.parse(data.Messages[0].Body);
            //make call to nodejs handler in codeigniter
            exec('php '+ appConf.CI_FC_PATH +'/index.php nodejs_handler make_contentq_call "'+ sqs_message_body.contentq_cat_id+'" "'+sqs_message_body.cnhq_cat_id+'" "'+sqs_message_body.network_id+'"',
              function (error, stdout, stderr) {
                if (error) {
                  // error handling 
                }
                if(stdout == 'Success'){
                  //delete message from queue
                  sqs.deleteMessage({
                    "QueueUrl" : appConf.sqs_distribution_url,
                    "ReceiptHandle" :data.Messages[0].ReceiptHandle
                  }, function(err, data){                
                  });
                }
                readMessage();                
              });
          }          
        }        
        readMessage();        
      });
    }
    readMessage();
    

    关于内存泄漏:我认为你不应该担心,因为readMessage() 的下一次调用发生在回调函数中 - 所以不是递归的,递归调用的函数在调用 receiveMessage() 函数后将值返回给父函数。

    【讨论】:

    • 您好,请查看我的这个要点。 gist.github.com/yalamber/374add88e887e688d818
    • 另外我是否应该担心运行此脚本时的内存泄漏?
    • 旧帖子,但是您是否遇到过内存泄漏的任何问题,因为我使用与 Bluebird 承诺几乎相同的想法。在每次轮询中,我都可以看到堆的总数增长了大约 200kb。
    • 我认为这种方法仍然会泄漏内存。
    • @napalm 请添加更多详细信息,这会有所帮助
    猜你喜欢
    • 2018-03-05
    • 2012-06-23
    • 2016-10-07
    • 2020-08-22
    • 2018-09-10
    • 2020-02-14
    • 2015-04-23
    • 2021-10-27
    • 2018-10-25
    相关资源
    最近更新 更多