【问题标题】:Storm Word Count Topology - Concept issue with number of executionsStorm Word Count Topology - 执行次数的概念问题
【发布时间】:2015-10-20 00:45:10
【问题描述】:

下午好,我正在关注 Storm-starter WordCountTopology here。作为参考,这里是 Java 文件。

这是主文件:

public class WordCountTopology {
public static class SplitSentence extends ShellBolt implements IRichBolt {

public SplitSentence() {
  super("python", "splitsentence.py");
}

@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
  declarer.declare(new Fields("word"));
}

@Override
public Map<String, Object> getComponentConfiguration() {
  return null;
}
}

public static class WordCount extends BaseBasicBolt {
Map<String, Integer> counts = new HashMap<String, Integer>();

@Override
public void execute(Tuple tuple, BasicOutputCollector collector) {
  String word = tuple.getString(0);
  Integer count = counts.get(word);
  if (count == null)
    count = 0;
  count++;
  counts.put(word, count);
  collector.emit(new Values(word, count));
}

@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
  declarer.declare(new Fields("word", "count"));
}
}

public static void main(String[] args) throws Exception {

TopologyBuilder builder = new TopologyBuilder();

builder.setSpout("spout", new TextFileSpout(), 5);

builder.setBolt("split", new SplitSentence(), 8).shuffleGrouping("spout");
builder.setBolt("count", new WordCount(), 12).fieldsGrouping("split", new Fields("word"));

Config conf = new Config();
conf.setDebug(true);

if (args != null && args.length > 0) {
  conf.setNumWorkers(3);

  StormSubmitter.submitTopology(args[0], conf, builder.createTopology());
}
else {
  conf.setMaxTaskParallelism(3);
  LocalCluster cluster = new LocalCluster();
  cluster.submitTopology("word-count", conf, builder.createTopology());
  Thread.sleep(10000);
  cluster.shutdown();
}
}
}

我不想从一个随机的 String[] 中读取,而是从一个句子中读取一个:

public class TextFileSpout extends BaseRichSpout {
    SpoutOutputCollector _collector;
    String sentence = "";
    String line = "";
    String splitBy = ",";
    BufferedReader br = null;

    @Override
    public void open(Map conf, TopologyContext context,
            SpoutOutputCollector collector) {
        _collector = collector;

    }

    @Override
    public void nextTuple() {
        Utils.sleep(100);
        sentence = "wordOne wordTwo";
        _collector.emit(new Values(sentence));
        System.out.println(sentence);
    }

    @Override
    public void ack(Object id) {
    }

    @Override
    public void fail(Object id) {
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("word"));
    }

}

这段代码运行并且输出是很多线程/发射。问题是程序执行重复读取一个句子 85 次而不是一次。我猜这是因为原始代码多次执行新的随机句子。

是什么导致 NextTuple 被调用这么多次?

【问题讨论】:

  • 能分享一下你的spout代码吗
  • @user2720864 共享喷口代码。对此感到抱歉

标签: java apache-storm word-count


【解决方案1】:

您应该使用 open 方法移动文件初始化代码,否则每次调用 nextTuple 时,您的文件处理程序都会被初始化。

编辑:

在 open 方法中,做类似的事情

    br = new BufferedReader(new FileReader(csvFileToRead));

然后读取文件的逻辑应该在 nextTuple 方法中

     while ((line = br.readLine()) != null) {
         // your logic
     }

【讨论】:

  • 我已将文件初始化移至打开状态。生成的句子用空格分隔文件中的所有单词进行更正。但是,nextTuple 被调用了 86 次,所以我的计数是应有的 86 倍。我想这将我的整个问题缩小到弄清楚如何只调用一次 nextTuple 。非常感谢您的宝贵时间。
  • 阅读逻辑应该在 nextTuple 方法中,更新我的答案
  • 感谢您的回复。即使我删除了整个文件读取部分,只将句子变成一个单词,nextTuple 也会重复调用 85 次。你知道吗,在这个例子中,Storm 如何决定运行 nextTuple 多少次?也许我错过了某处的配置。谢谢。
  • 我将代码简化为只读取一个句子。我查看了代码,不知道什么叫 NextTuple。我需要 spout 运行一次并返回 wordOne: 1 和 wordTwo: 1 而不是 85 和 85。谢谢
  • Storm 专为在数据可用时发出数据的流式源而设计。 nextTuple() 在无限循环中被调用,因此对于您的情况,它需要跟踪它在数据源中的位置。如果您希望至少处理一次,它还应该跟踪 ack() 和 fail() 调用。
猜你喜欢
  • 2011-10-22
  • 1970-01-01
  • 1970-01-01
  • 2011-03-11
  • 1970-01-01
  • 1970-01-01
  • 2011-07-31
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多