在Storm上使用BaseRichSpout调用nextTuple()无限次

ray*_*man 4 java bigdata apache-storm

我实现了简单的Storm拓扑,它具有单个喷口和在本地集群模式下运行的螺栓.

由于某种原因,spout的nextTuple()不止一次被调用.

知道为什么吗?

码:

喷口:

public class CommitFeedListener extends BaseRichSpout {
    private SpoutOutputCollector outputCollector;
    private List<String> commits;

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

    @Override
    public void open(Map configMap,
                     TopologyContext context,
                     SpoutOutputCollector outputCollector) {
        this.outputCollector = outputCollector;
    }

    **//that method is invoked more than once**
    @Override
    public void nextTuple() {

            outputCollector.emit(new Values("testValue"));

    }
}
Run Code Online (Sandbox Code Playgroud)

螺栓:

public class EmailExtractor extends BaseBasicBolt {
    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("email"));
    }
    @Override
    public void execute(Tuple tuple,
                        BasicOutputCollector outputCollector) {
        String commit = tuple.getStringByField("commit");
        System.out.println(commit);        
    }  
}
Run Code Online (Sandbox Code Playgroud)

运行配置:

public class LocalTopologyRunner {
    private static final int TEN_MINUTES = 600000;
    public static void main(String[] args) {
        TopologyBuilder builder = new TopologyBuilder();
        builder.setSpout("commit-feed-listener", new CommitFeedListener());
                builder
        .setBolt("email-extractor", new EmailExtractor())
                .shuffleGrouping("commit-feed-listener");
        Config config = new Config();
        config.setDebug(true);
        StormTopology topology = builder.createTopology();
        LocalCluster cluster = new LocalCluster();
        cluster.submitTopology("github-commit-count-topology",
                config,
                topology);
        Utils.sleep(TEN_MINUTES);
        cluster.killTopology("github-commit-count");
        cluster.shutdown();
    }
}
Run Code Online (Sandbox Code Playgroud)

雷万,谢谢大家.

zen*_*eni 6

nextTuple()在设计中以无限循环方式调用.像这样使用例如对外部资源(数据库,流,IO等)的脏检查.

如果你在nextTuple()中没有任何关系,你应该睡一会儿以防止使用backtype.storm.utils.Utils发送CPU垃圾邮件

Utils.sleep(pollIntervalInMilliseconds);
Run Code Online (Sandbox Code Playgroud)

Storm是一种实时处理架构,因此它确实是正确的行为.检查一些样品,看看如何根据您的需要实施喷口.