【发布时间】:2017-02-26 12:45:26
【问题描述】:
一段时间以来,我一直在尝试了解如何创建一组订阅单个主题的 kafkaspout。我发现一些引用 spoutconfig 中的 id 的来源可以用作组 id。我的主要问题是我如何知道创建的 spout 是否作为一个组起作用。我还想知道setSpout() 中的 paralellism_hint 是否是创建一组 paralellism_hint 数量的 spout 的人。请赐教。
我的代码是这样的
private KafkaSpout buildKafkaSpout(String zkTopic, String zkRoot, String groupId) {
String zkConnString = "localhost:2181";
BrokerHosts hosts = new ZkHosts(zkConnString);
SpoutConfig kafkaSpoutConfig = new SpoutConfig (hosts, zkTopic,zkRoot, groupId );
kafkaSpoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());
KafkaSpout kafkaSpout = new KafkaSpout(kafkaSpoutConfig);
return kafkaSpout;
}
private void buildTopologyForGroupOne(TopologyBuilder builder ) {
String zkTopic = "topic1";
String zkRoot = "/topic1";
String groupId = "group1";
List<String> zkSpoutIds = new ArrayList<String>();
zkSpoutIds.add("word_count-spout");
zkSpoutIds.add("total_word_count-spout");
for(String spoutId:zkSpoutIds ){
KafkaSpout kafkaSpout = buildKafkaSpout(zkTopic, zkRoot, groupId);
builder.setSpout(spoutId.concat("_"+groupId), kafkaSpout,2);
}
builder.setBolt("word_split-bolt_group_1",new SplitBolt()).shuffleGrouping("word_count-spout"+"_"+groupId);
builder.setBolt("split_count-bolt_group_1",new CountBolt()).shuffleGrouping("word_split-bolt_group_1");
builder.setBolt("total_word_count-bolt_group_1",new TotalWordCountBolt()).shuffleGrouping("total_word_count-spout"+"_"+groupId);
}
现在我如何知道我创建的两个 spout(word_count-spout,total_word_count-spout) 是否作为一个组。作为一个组,我的意思是如果创建了一个新的 spout,zookeeper 将重新排列分区。
提前致谢
【问题讨论】:
-
终于找到了来源。可以使用新版本的storm apache-storm-1.0.2 创建一个带有consumer id 的KafkaSpout。
标签: apache-kafka apache-storm kafka-consumer-api kafka-producer-api