【问题标题】:How do I output a stream of tuples from a Storm spout with emit() and sync()?如何使用 emit() 和 sync() 从 Storm spout 输出元组流?
【发布时间】:2015-01-25 04:52:19
【问题描述】:

(xpost github issue)

我是 Storm 的新手。我找到了有用的 node-storm 库,并且我已经成功提交了拓扑,但我无法让我的 spout 发出元组流。

node-storm 的 wordcount 示例运行良好。

我想要一个订阅 websocket 并将任何消息作为元组输出的 spout。

到目前为止,这是我的尝试。我想我有一些错误配置,因为我知道我的 wsEmitter 正在发出 future 事件,但我的 Storm UI 显示零喷口发射。

我怀疑也许我不应该在 spout 函数中绑定监听器?

这个函数会被多次调用吗? (看起来像......见https://github.com/RallySoftware/node-storm/blob/master/lib/spout.js#L4

sync 的实际作用是什么?我应该什么时候使用它?

var storm = require('node-storm');
var wsEmitter = require('./wsEmitter.js')();
wsEmitter.init();  // subscribe to websocket

var futuresSpout = storm.spout(function(sync) {
  var self = this;
  console.log('subscribing to ws');
  wsEmitter.on('future', function(data){       // websocket data arrived
    self.emit([data]);
    sync();
  });
})
.declareOutputFields(["a"]);

【问题讨论】:

    标签: node.js apache-storm hortonworks-data-platform


    【解决方案1】:

    原来我有两个问题。首先,我的拓扑没有执行,因为我的一个螺栓(未显示)未能设置.declareOutputFields()

    其次,我需要延迟 spout 的发射,直到主管用 nextTick() 要求发射。我通过缓冲任何传入的消息来做到这一点,直到主管调用 spout:

    module.exports = (function(){
      var storm = require('node-storm');
    
      var wsEmitter = require('./wsEmitter.js')();
      wsEmitter.init();
    
      var queue = [];
      var queueEmpty = true;
    
      wsEmitter.on('thing', function(data){
        var trade = JSON.parse(data);
        trade.timeReported = new Date().valueOf();
        queue.push(trade);
        queueEmpty = false;
      });
    
      return storm.spout(function(sync) {
        var self = this;
        setTimeout(function(){
          if(!queueEmpty){
            self.emit([queue.shift()]);
            queueEmpty =
            ( queue.length === 0
            ? true
            : false )
          }
          sync();
        }, 100);
    
      })
      .declareOutputFields(['trade'])
    })()
    

    【讨论】:

      猜你喜欢
      • 2018-08-14
      • 1970-01-01
      • 2013-05-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-12-30
      • 1970-01-01
      • 2015-04-22
      相关资源
      最近更新 更多