我在我的机器和我的其他机器上安装了Storm 1.1.1,我使用的是Kafka版本0.10.0.1.这两个服务都与Zookeeper 3.4.6版连接
我成功部署了我的拓扑,看起来像这样:
public class SOTopology
{
public static void main (String[] args ) throws Exception
{
final String brokers = args[0];
final String kafkaTopic = args[1];
final String mongo_uri = args[2];
final String mongo_collection = args[3];
TopologyBuilder topology=new TopologyBuilder();
topology.setSpout("KafkaSpout",new KafkaSpout<>(KafkaSpoutConfig.builder(brokers, kafkaTopic).build()), 1);
topology.setBolt("FilterBolt", new Filterbolt(),1).shuffleGrouping("KafkaSpout");
topology.setBolt("TagCountBolt", new TagCountBolt(),1).shuffleGrouping("FilterBolt");
topology.setBolt("TopicBolt", new TopicBolt(),1).shuffleGrouping("FilterBolt");
topology.setBolt("MongoDBBolt",new MongoDBBolt(),1).shuffleGrouping("TagCountBolt").shuffleGrouping("TopicBolt");
Config conf = new Config();
conf.setDebug(true);
conf.put("mongo.uri", mongo_uri);
conf.put("mongo.collection", mongo_collection);
conf.setMaxSpoutPending(40);
conf.setNumWorkers(10);
conf.setDebug(true);
StormSubmitter.submitTopology("StackOverflowTopology", conf, topology.createTopology());
}
}
Run Code Online (Sandbox Code Playgroud)
当我去我的StormUI时,我收到以下消息:Offset lags for kafka not …