【问题标题】:How to synchronize separate callbacks in Javascript如何在Javascript中同步单独的回调
【发布时间】:2017-02-16 20:51:33
【问题描述】:

我正在开发一个 nodejs 应用程序(使用 express),它将通过 HTTP 为来自客户端的 REST 调用提供服务。其中一个 REST API 将处理 POST 请求,该请求将从 POST 正文获取数据并通过 MQTT 客户端(作为应用程序的一部分启动)发布。

订阅主题的单独应用程序(通过我们都连接到的 MQTT 代理)将接收消息并通过向我的应用程序的 MQTT 客户端订阅的主题发布消息来响应(这将导致在我的 Javascript 应用中触发的回调)。

我希望能够返回我的应用程序在我的应用程序正在处理的 POST 响应中收到的 MQTT 消息。

在常规的“C 语言编程”术语中...有一个线程处理 REST API 和一个单独的线程处理接收 MQTT 消息(例如两个不同的套接字)。我想阻止线程处理信号量上的 POST 处理线程,直到 MQTT 线程可以接收一些数据,将其入队并释放信号量以解除对 POST 处理线程的阻塞(然后将消息出列并在POST 响应)。

在修改了各种模块和 Promise/Generators 之后,我不清楚如何使它工作......

执行此操作的“javascript 方式”是什么?

TIA!

【问题讨论】:

  • 关于代码的问题应该包括您的实际代码。

标签: javascript node.js rest express mqtt


【解决方案1】:

您不会在 node.js 中阻塞。根据它的设计,node.js 是事件驱动的,所有网络 I/O(以及大多数文件 I/O)都是非阻塞和异步的。因此,您需要协调多个异步操作并在两者都完成时收到通知(通过某种事件)。

当您不共享任何代码时,很难非常具体,但如今协调异步操作的主要工具是 Promise。因此,如果您让每个异步操作都返回一个承诺,该承诺将在异步操作完成时解决或在异步操作遇到错误时拒绝,那么您可以将Promise.all() 与两个承诺一起使用,它会告诉您何时异步操作已完成并为您提供两种结果。

大致思路是这样的:

 let p1 = asyncOperation1(...);
 let p2 = asyncOperation2(...);

 Promise.all([p1, p2]).then(results => {
     // both async operations are done here
     // results[0] and results[1] contain the resolved value of each of the two promises
 }).catch(err => {
     // process error here
 });

为了从 REST 操作中获取 Promise,有许多不同的解决方案,从手动围绕操作执行您自己的 Promise 包装器,到包装请求模块以向具有 .promisify() 和 @ 的 Bluebird 等库提供 Promise 的模块987654325@ 自动为您包装操作的方法。要提供有关如何从异步操作中获得 Promise 的更具体信息,您必须分享实际代码。


由于您似乎是新来的,让我提供一些发布建议。如果您向我们展示您的代码,我们将始终能够提供更完整和更相关的答案。当您不向我们展示您的代码时,您实际上是在要求一个通用教程,它更难编写,也更难确保它涵盖您的确切用例。

【讨论】:

    【解决方案2】:

    不要阻止,在您发布的消息中包含唯一 ID,并让响应消息也包含此 ID。

    然后,您将快速响应对象存储在使用此 ID 作为键的对象中,以便当响应到达您并发回响应时。

    您还应该包含一个时间戳,这样您就可以定期遍历等待的响应对象,并在 MQTT 消息未及时到达时做出响应。

    var onGoingRequests = {};
    var id = 0;
    
    mqttClient.on('message',function(topic,message){
      var payload = JSON.parse(message.toString());
      var details = onGoingRequest[payload.id];
      if (details) {
        details.response.status(200).send(details.body);
        delete onGoingRequests[payload.id];
      } else {
        //response too late
      }
    });
    
    app.post('/foo', function(req, resp) {
        var message = {
          id: 'foo' +id++,
          body: req.body
        }
        var topic = 'request/foo'; 
        mqttClient.publish(topic, JSON.stringify(message));
        onGoingRequests[message.id] = {
          response: resp,
          timestamp: Date.now(),
        }; 
    });
    
    var timeout = setInterval(function(){
      var now = Date.now();
      var keys = Object.keys(onGoingCommands);
      for (key in keys){
        var waiting = onGoingCommands[keys[key]];
        if (waiting) {
          var diff = now - waiting.timestamp;
          if (diff < timeout) {
            waiting.res.status(504).send('{"error": "timeout"}');
            delete onGoingCommands[keys[key]];
          }
        }
      }
    }, 500);
    

    【讨论】:

      【解决方案3】:

      最终发现我可能一直在考虑解决方案......

      在 POST 处理回调中,我将响应对象排入队列并返回。稍后,在 MQTT 订阅消息处理回调中,我将挂起的响应对象出列并使用它向挂起的 POST 事务发送正确的响应。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2012-01-02
        • 2019-06-08
        • 1970-01-01
        • 2013-01-13
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多