【问题标题】:Processing a Message Queue and Using Async处理消息队列并使用异步
【发布时间】:2023-03-24 09:22:01
【问题描述】:

我编写了一个小型测试节点应用程序,它循环并将消息添加到队列(天蓝色存储队列),如下所示:

var queueService = azure.createQueueService();
var queueName = 'taskqueue';

// other stuff like check if created

// loop called after queue is confirmed

for (i=0;i<1000;i++){
  queueService.createMessage(queueName, "Hello world!", null, messageCreated);
}

// messageCreated does nothing at the moment, just logs to console

我正在尝试重写它来处理说 100 万个使用异步来控制并行运行的工作函数的数量的创建。这比什么都重要。

https://github.com/caolan/async#queue

这是异步队列的基本设置,我对需要更改的内容有点茫然。我认为以下方法行不通:

var q = async.queue(function (task, callback) {
   queueService.createMessage(queueName, task.msg, null, messageCreated);
    callback();
}, 100);


// assign a callback.  Called when all the queues have been processed
q.drain = function() {
    console.log('all items have been processed');
}

// add some items to the queue
for(i=0;i<1000000;i++) {
  q.push({msg: 'Hello World'}, function (err) {
    console.log('finished processing foo');
  });
    console.log('pushing: ' + i);
}

我不太明白如何将它们与异步结合在一起。

【问题讨论】:

  • 快速浏览一下我觉得没问题..?
  • @Alfred,我将百万改为 500 进行测试。最终发生的事情是我看到“完成处理 foo”打印到屏幕 500 次,然后是 messageCreated 的第一个控制台日志(来自 azure put 消息的回调)..

标签: node.js asynchronous azure azure-storage


【解决方案1】:

这是你的错误:

var q = async.queue(function (task, callback) {
   queueService.createMessage(queueName, task.msg, null, messageCreated);
    callback();
}, 100);

您正在做的是在队列上创建一条消息,然后立即传递延续(调用回调)。您要做的是在传递给createMessage 的回调中传递延续:

var q = async.queue(function (task, callback) {
   queueService.createMessage(queueName, task.msg, null, function(error, serverQueue, serverResponse) {
       callback(error, serverQueue, serverResponse);
       messageCreated(error, serverQueue, serverResponse);
   });
}, 100);

现在每个任务都会在实际创建任务后报告完成。

编辑:更新了createMessage 回调的接口。

【讨论】:

  • createMessagecallbackfunction messageCreated(error, serverQueue, serverResponse),所以我不能像你这样称呼它。这就是我的困惑所在。
  • 在进行推送的 for 循环中,如果我在推送后使用 console.log:console.log('pushing: ' + i); 你会希望看到i 打印出 1000000 次然后队列开始处理还是会您希望队列在处理 for 循环的同时开始处理?例如,我希望看到pushing: 1 pushing: 2 pushing: 3 finished processing foo pushing: 4
  • 通过处理我的意思是 createMessage 的回调将在 for 循环完成之前触发。
  • 这很难说,取决于createMessagepush 的速度。你告诉我!
  • 我可以告诉你,运行一百万需要一两分钟,而且在 for 循环完成处理之前,我没有一次收到 createMessage 的回调。推送通常需要一两秒才能返回。
猜你喜欢
  • 1970-01-01
  • 2015-12-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-09-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多