【问题标题】:Producer-Consumer Queue in AngularJSAngularJS 中的生产者-消费者队列
【发布时间】:2015-11-04 18:50:11
【问题描述】:

几年前我就知道python和数据库了。

但我想提高我有限的 JavaScript 知识。对于我的玩具项目,我想在 Web 浏览器中使用异步队列并为此使用 AngularJS。

在 python 中有一个很好的类,叫做multiprocessing.Queue,我以前用过。

现在我搜索类似的东西,但在 AngularJS 中

  • 第 1 步:队列中拉出工作项(粉红色圆圈)。只是查看json字节。

  • 第 2 步:用户处理数据。

  • 第 3 步:out-queue 负责将结果发送到服务器。

为什么会有这种“复杂”的设置?因为我希望应用程序尽可能地响应。入队应预加载一些数据,出队应处理响应通信。

另一个好处是,通过这种设置,应用程序可以处理几分钟的服务器或网络中断。

AngularJS 的双向数据绑定会立即更新用户编辑的数据,这并不适合我的问题。或者我错过了什么。我是 AngularJS 的新手。

图中粉红色的圆圈代表 JSON 数据结构。我想通过一个请求将它们中的每一个推送到浏览器。

例子:

用户看到一个问题,然后他需要填写三个字段。例如:

  • 答案:输入文字
  • like-this-question: integer from 1..5
  • 难度:从 1..5 开始的整数

数据应该在按下“提交”后放入队列。他应该马上得到下一个问题。

问题:

AngularJS 是否已有可用的生产者-消费者队列?如果没有,如何实现?

更新

从客户端发送数据可以使用纯 AJAX 实现。预取数据的队列是更复杂的部分。尽管两者都可以使用相同的实现。客户端以超低延迟获取新数据非常重要。 in-queue 每次最多填充 5 个 item 以避免客户端等待数据。

在我的情况下,浏览器是否关闭并且队列中的项目丢失并不重要。填充队列在服务器部分是只读的。

我没有固定在 AngularJS 上。如果有充分的理由,我很乐意更改框架。

可以使用 localStorage (html5) 来保留浏览器重新加载之间的队列

【问题讨论】:

  • 您是否需要在 JS 中实现队列,或者这只是您的偏好?为什么不直接使用现有的多种队列之一并通过 REST 与之通信?我的理解是AngularJS应该以这种方式使用。
  • npmjs.com/package/node-taskqueue 如果您有兴趣使用用 JS 编写的任务队列 API。

标签: javascript angularjs producer-consumer


【解决方案1】:

重新思考一下,你的前端真的需要生产者-消费者吗? 在您的示例中,我认为在这种情况下,一个简单的 pub-sub$Q 就足够了。

示例:

创建一个订阅者服务,在您的示例中,questionHandlerService 订阅事件question-submitted,在您的questionService 中,当用户提交问题时,只需使用数据发布事件question-submitted,触发并忘记。您无需等待questionHandlerService 的回复。

请记住,javascript 中只有一个主线程,如果您的方法会阻塞 ui,例如 loop through 1000 items in the array, synchronous process them,如果您将其放在另一个“队列”中将无济于事,因为它必须执行时阻止 ui,除非您在 web-worker 中执行它。如果用户刷新浏览器怎么办?您未处理的请求刚刚丢失。

触发 XHR 调用不会阻塞 UI,在前端实现 producer-consumer 没有意义,您只需验证输入触发 XHR,让后端处理繁重的工作,在您的后端控制器中,您可以使用队列来立即保存请求和响应。并从任何线程或其他进程处理队列。

【讨论】:

  • 是的,你是对的。从客户端发送数据可以使用我们的解决方案或简单的 AJAX 来实现。预取数据的队列是更复杂的部分。客户端以超低延迟获取新数据非常重要。 in-queue 每次最多填充 5 个 item 以避免客户端等待数据。
【解决方案2】:

Working Plunker - 后端服务仅用于模拟其余接口;我修复了一些错误,例如错误限制。所以,假设 Plunker 是最后一个版本...

我没有足够的时间来改进以下内容,但是,您可以将其作为起点...

顺便说一句,我认为你需要的可能是:

  1. 包装 $http 的服务。
  2. 我们需要注册任务时使用的方法 push
  3. 递归私有方法_each,逐步减少注册队列。
  4. ...您认为重要的所有其他内容(getCurrentTask、removeTask 等)。

用法:

angular
  .module('myApp', ['Queue'])
  .controller('MyAppCtrl', function($httpQueue, $scope) { 
    var vm = $scope;
    
    // using a route.resolve could be better!
    $httpQueue
      .pull()
      .then(function(tasks) { vm.tasks = tasks;  })
      .catch(function() { vm.tasks = [{ name: '', description: '' }]; })
    ;
  
    vm.onTaskEdited = function(event, task, form) {
      event.preventDefault();
      if(form.$invalid || form.$pristine ) { return; }
      
      
      return $httpQueue.push(task);
      
    };
  })
;
<article ng-app="myApp">
  <div ng-controller="MyAppCtrl">
    
    
    <form ng-repeat="task in tasks" name="taskForm" ng-submit="onTaskEdited($event, task, taskForm)">
      <input ng-model="task.name" placeholder="Task Name" />
      <textarea ng-model="task.description"></textarea>
    </form>
    
    
  </div>
</article>

模块定义

(function(window, angular, APP) {
  'use strict';

  function $httpQueueFactory($q, $http) {
    var self = this;

    var api = '/api/v1/tasks';

    self.queue = [];
    var processing = false;


    //Assume it as a private method, never call it directly
    self._each = function() {
      var configs = { cache: false };
			
			
      return self
        .isQueueEmpty()
        .then(function(count) {
          processing = false;
          return count;
        })
        .catch(function() {
          if(processing) {
            return;
          }
        
          processing = true;
          var payload = self.queue.shift();
     
          var route = api;
          var task = 'post';
          if(payload.id) {
            task = 'put';
            route = api + '/' + payload.id;
          }
        
          return $http
            [task](route, payload, configs)
            .catch(function(error) {
              console.error('$httpQueue._each:error', error, payload);
              //because of the error we re-append this task to the queue;
              return self.push(payload);
            })
            .finally(function() {
              processing = false;
              return self._each();
            })
          ;
        })
      ;
    };


    self.isQueueEmpty = function() {
      var length = self.queue.length;
      var task = length > 0 ? 'reject' : 'when';
      
      return $q[task](length);
    };

    self.push = function(data) {
      self.queue.push(data);
      self._each();

      return self;
    };

    self.pull = function(params) {
      var configs = { cache: false };
      
      configs.params = angular.extend({}, params || {});

      return $http
        .get(api, configs)
        .then(function(result) {
          console.info('$httpQueue.pull:success', result);

          return result.data;
        })
        .catch(function(error) {
          console.error('$httpQueue.pull:error', error);
        
          return $q.reject(error);
        })
      ;
    };
  }



  APP
    .service('$httpQueue', ['$q', '$http', $httpQueueFactory])
  ;	

})(window, window.angular, window.angular.module('Queue', []));

处理 DataLayer 更改

处理数据层上的更改(您称为队列中的)反而是一项更困难的任务,因为我们需要保持 最后一次拉动 与当前拉动...

顺便说一句,如果您正在寻找具有实时通知的系统,您可能应该看看 套接字层...我建议 Socket.io 因为它是经过充分测试、行业认可的解决方案。

如果您无法实现套接字层,另一种解决方案可能是长轮询实现,详细描述here

为简单起见,在这篇文章中,我们将实现一个简单的间隔来更新当前任务列表... 所以,前面的例子变成了:

angular
  .module('myApp', ['Queue'])
  .controller('MyAppCtrl', function($httpQueue, $scope, $interval) { 
    var vm = $scope;
    var 
      pollingCount = 0, // infinite polling
      pollingDelay = 1000
    ;
    
    // using a route.resolve could be better!
    $httpQueue
      .pull()
      .then(function(tasks) { vm.tasks = tasks;  })
      .catch(function() { vm.tasks = [{ name: '', description: '' }]; })
      .finally(function() { return $interval(vm.updateViewModel.bind(vm), pollingDelay, pollingCount, true); })
    ;
  
    var isLastPullFinished = false;
    vm.updateViewModel = function() {
      if(!isLastPullFinished) { return; }
      
      return $http
        .pull()
        .then(function(tasks) {
          for(var i = 0, len = tasks.length; i < len; i++) {
            
            for(var j = 0, jLen = vm.tasks.length; j < jLen; j++) {
              if(tasks[i].id !== vm.tasks[j].id) { continue; }
              
              // todo: manage recursively merging, in angular 1.3+ there is a
              // merge method https://docs.angularjs.org/api/ng/function/angular.merge
              // todo: control if the task model is $dirty (if the user is editing it)
              angular.extend(vm.tasks[j], tasks[i]);
            }
            
          };
          
          return vm.tasks;
        })
        .finally(function() {
          isLastPullfinished = true;
        })
      ;
    };
  
    
    
    vm.onTaskEdited = function(event, task, form) {
      event.preventDefault();
      if(form.$invalid || form.$pristine ) { return; }
      
      
      return $httpQueue.push(task);
      
    };
  })
;

希望对您有所帮助; 我没有测试过,所以,可能有一些错误!

【讨论】:

  • 我还没有在 JS 中使用递归方法(仅在其他语言中)。不知何故,我认为这只是一个无限循环。你为什么对_each()使用递归方法调用?
  • 感谢您提供指向 socket.io 的链接。看起来坚如磐石。谢谢
  • 它不是无穷无尽的,当队列数组的长度等于 0 时,promise 将被解析,_each 方法将不再调用......方式,我认为你应该我的解决方案......可能对你有用!
  • 对不起,我不明白你写的这句话:“顺便说一句,我认为你应该是我的解决方案。”。你想告诉我什么?
  • 对不起,我错过了尝试,我的英语很差,但我正在学习!您应该尝试我在下面发布的代码! :)
【解决方案3】:

我会研究 Promise,因为它们提供了您似乎需要的异步功能,可能还结合了在初始化时引入所有或一些问题。如果您确实需要离线查看服务工作者,尽管它们仅限于较新的浏览器

如果您不能使用服务人员,您可以创建一个服务(或工厂)来保存您的所有问题或提前阅读。您将使用 $http 服务来尝试获取更多答案,如果在网络连接再次出现之前您无法使用您所拥有的东西。

例如,通过 $interval 尝试在循环中获取新答案(并发布答案),可以使用 $http 服务来检查网络。

【讨论】:

  • 服务人员太新了。很遗憾,浏览器支持还不在这里:caniuse.com/#feat=serviceworkers。我想我需要一个单一的页面视图。这样我就可以使用 AppCache 或在自定义 JS 中保持入队和出队。
【解决方案4】:

应用结构: 2 个数组,一个带问题,另一个带答案。

应用程序启动: 两个 $interval 对象,一个从服务器获取前 10 个问题。如果问题缓冲区长度小于 10,则会将新问题推送到数组中。

另一个函数检查答案数组,如果通信可用,则将结果发送到服务器。

当用户检查一个问题时,它会从第一个数组中弹出并推送到答案数组中。

就是这样......数据绑定的 Angular 2 方式使用调用应用程序中的函数(如 giveMeTheLatestQuestion...)显示最古老的问题......

希望对你有帮助!

【讨论】:

  • 我不明白这样做的好处,如果用户在此期间刷新页面怎么办?您可能只处理了其中的 5 个。另外5个刚刚输了?在前端实现这样的拱门是没有意义的
  • @Sean 我更新了问题:在我的情况下,浏览器是否关闭并且队列中的项目丢失并不重要。填充队列在服务器部分是只读的。
  • @MatteoConta 谢谢你的回答。您的回答描述了一种我可以用来实现的算法。我的想法真的那么奇怪,没有可重复使用的东西吗?
  • @Sean 使用 2 个间隔较短的计时器可能会在用户刷新页面时限制问题。在我看来,整个页面重新加载对于每个架构都是一个问题,您需要实现服务器端会话以限制此类问题。 Guettli 感谢您的回答,我会从头开始编写代码,我不知道库或框架是否可以帮助您或使解决方案复杂化。
  • @guettli,看看这篇文章,可能对你的项目有用thinkster.io/a-better-way-to-learn-angularjs/promises
猜你喜欢
  • 2012-01-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多