Created
December 28, 2014 16:30
-
-
Save mdobson/0aa88c08d1b15e4b7203 to your computer and use it in GitHub Desktop.
Kafka Topology
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package storm.starter; | |
| import java.util.Map; | |
| import java.util.UUID; | |
| import storm.kafka.KafkaConfig.BrokerHosts; | |
| import storm.kafka.KafkaConfig.ZkHosts; | |
| import storm.kafka.KafkaSpout; | |
| import storm.kafka.SpoutConfig; | |
| import backtype.storm.Config; | |
| import backtype.storm.LocalCluster; | |
| import backtype.storm.StormSubmitter; | |
| import backtype.storm.task.OutputCollector; | |
| import backtype.storm.task.TopologyContext; | |
| import backtype.storm.topology.OutputFieldsDeclarer; | |
| import backtype.storm.topology.TopologyBuilder; | |
| import backtype.storm.topology.base.BaseRichBolt; | |
| import backtype.storm.tuple.Tuple; | |
| import backtype.storm.utils.Utils; | |
| import backtype.storm.spout.SchemeAsMultiScheme; | |
| import storm.kafka.StringScheme; | |
| /* | |
| * Simple Kafka Topology. Reads from a Kafka Queue, and sends to a printer bolt. | |
| * | |
| */ | |
| public class KafkaTopology { | |
| public static class PrinterBolt extends BaseRichBolt { | |
| OutputCollector _collector; | |
| @Override | |
| public void prepare(Map stormConf, TopologyContext context, | |
| OutputCollector collector) { | |
| // TODO Auto-generated method stub | |
| _collector = collector; | |
| } | |
| @Override | |
| public void execute(Tuple input) { | |
| // TODO Auto-generated method stub | |
| System.out.println(input.getString(0)); | |
| } | |
| @Override | |
| public void declareOutputFields(OutputFieldsDeclarer declarer) { | |
| // TODO Auto-generated method stub | |
| } | |
| } | |
| public static void main(String[] args) throws Exception { | |
| BrokerHosts hosts = new ZkHosts("localhost:2181", "/brokers"); | |
| SpoutConfig spoutConf = new SpoutConfig(hosts, "test", "/test", UUID.randomUUID().toString()); | |
| spoutConf.scheme = new SchemeAsMultiScheme(new StringScheme()); | |
| KafkaSpout kafkaSpout = new KafkaSpout(spoutConf); | |
| TopologyBuilder builder = new TopologyBuilder(); | |
| builder.setSpout("kafka", kafkaSpout, 10); | |
| builder.setBolt("printer", new PrinterBolt(), 3).shuffleGrouping("kafka"); | |
| Config conf = new Config(); | |
| conf.setDebug(true); | |
| if(args != null && args.length > 0) { | |
| conf.setNumWorkers(3); | |
| StormSubmitter.submitTopologyWithProgressBar(args[0], conf, builder.createTopology()); | |
| } else { | |
| LocalCluster cluster = new LocalCluster(); | |
| cluster.submitTopology("test", conf, builder.createTopology()); | |
| Utils.sleep(10000); | |
| cluster.killTopology("test"); | |
| cluster.shutdown(); | |
| } | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment