【发布时间】:2015-07-30 16:54:00
【问题描述】:
我从互联网上借用了代码来创建一个从远程 kafka 集群读取数据的storm spout。我检查了storm集群和kafka集群之间的连接,没问题。我可以通过 kafka 命令行工具阅读主题。但是当我提交拓扑时,它不会发出任何东西。请帮忙!!
public class KafkaTopology {
public static void main(String[] args) throws Exception {
List<String> hosts = new ArrayList<String>();
hosts.add("10.87.36.80:2181");
// SpoutConfig kafkaConf = new SpoutConfig(StaticHosts.fromHostString(hosts, 3), "test", "/kafkastorm", "discovery");
SpoutConfig kafkaConf = new SpoutConfig(new ZkHosts("localhost:2181"),
"test", // Kafka topic to read from
"/test", // Root path in Zookeeper for the spout to store consumer offsets
"clickdata"); // ID for storing consumer offsets in Zookeeper
kafkaConf.scheme = new SchemeAsMultiScheme(new StringScheme());
KafkaSpout kafkaSpout = new KafkaSpout(kafkaConf);
TopologyBuilder builder = new TopologyBuilder();
builder.setSpout("spout", kafkaSpout, 2);
builder.setBolt("printer", new PrinterBolt())
.shuffleGrouping("spout");
Config config = new Config();
config.setDebug(true);
if(args!=null && args.length > 0) {
config.setNumWorkers(3);
StormSubmitter.submitTopology(args[0], config, builder.createTopology());
} else {
config.setMaxTaskParallelism(3);
LocalCluster cluster = new LocalCluster();
cluster.submitTopology("kafka", config, builder.createTopology());
Thread.sleep(10000);
cluster.shutdown();
}
}
}
【问题讨论】:
-
你如何观察你的拓扑结构?使用风暴用户界面?您是否检查了日志中的错误消息?如何开始拓扑(在某些 IDE 中或使用
storm jar命令)? -
我也在观察工作节点和用户界面的日志。还使用 jar 命令提交了拓扑。
-
我你测试了从storm集群到kafka集群的连接,你使用了哪个storm集群节点(nimbus,all)?我也想知道
new ZkHosts("localhost:2181")是否可能是问题所在。你的 ZK 设置是什么? -
我正在使用 hortonworks 沙箱。 localhost对应zookeeper。日志还显示我已连接到 kafka 集群。我会粘贴日志。
-
你有本地运行的zookeeper吗?你的卡夫卡也是本地的呢?你的 kafkaspout 看起来不错...
标签: java apache-kafka apache-storm