【发布时间】:2015-08-06 21:04:16
【问题描述】:
我正在尝试将 Kafka 与 Storm 集成。我正在使用 Kafka Spout 从 Kafka 主题中检索数据并将其提供给 Storm Bolt 以进行进一步处理。我能够成功提交拓扑,但 Spout 没有发出任何信息data.It 也不会引发任何错误。我对 Kafka 和 Storm 很陌生。所以,我无法找到这个问题背后的原因。请提出修改建议。提前致谢!
我的拓扑:
public class TopologyMain {
private static final String SENTENCE_SPOUT_ID = "kafka-sentence-spout";
public static void main(String[] args) throws InterruptedException, AlreadyAliveException, InvalidTopologyException {
int numSpoutExecutors = 1;
TopologyBuilder builder = new TopologyBuilder();
builder.setSpout(SENTENCE_SPOUT_ID, buildKafkaSentenceSpout(), numSpoutExecutors);
builder.setBolt("word-normalizer", new WordNormalizer())
.shuffleGrouping(SENTENCE_SPOUT_ID);
builder.setBolt("word-counter", new WordCounter(),2)
.shuffleGrouping("word-normalizer");
//Configuration
Config conf = new Config();
conf.setDebug(false);
//Topology run
conf.put(Config.TOPOLOGY_MAX_SPOUT_PENDING, 1);
conf.put(Config.NIMBUS_HOST, "192.168.1.229");
conf.put(Config.NIMBUS_THRIFT_PORT, 6627);
System.setProperty("storm.jar", "/home/ubuntu/st/stIn/target/storm-wc.jar");
StormSubmitter.submitTopology("Count-Word-Topology", conf,builder.createTopology());
}
private static KafkaSpout buildKafkaSentenceSpout() {
BrokerHosts hosts = new ZkHosts("localhost:2181");
SpoutConfig spoutConfig = new SpoutConfig(hosts, "test", "/acking-kafka-sentence-spout", "acking-sentence-spout");
spoutConfig.forceFromStart = true;
spoutConfig.startOffsetTime = kafka.api.OffsetRequest.EarliestTime();
return new KafkaSpout(spoutConfig);
}
}
【问题讨论】:
-
您可以使用控制台脚本来消费主题吗?这个拓扑在本地工作吗?您可以从集群中的storm-starter 运行任何示例拓扑吗?
-
是的,我可以使用控制台脚本来处理主题。拓扑在本地运行良好。此外,我可以在 spout 是文本文件的集群上运行不同的风暴拓扑。
-
如果你的 zookeeper 运行在与
Config.NIMBUS_HOST相同的 ip 中,请尝试将 localhost:2181 更改为:2181 看看是否有帮助
标签: java apache-kafka apache-storm kafka-consumer-api