Skip to content

Instantly share code, notes, and snippets.

@mdobson
Created December 28, 2014 16:30
Show Gist options
  • Select an option

  • Save mdobson/0aa88c08d1b15e4b7203 to your computer and use it in GitHub Desktop.

Select an option

Save mdobson/0aa88c08d1b15e4b7203 to your computer and use it in GitHub Desktop.
Kafka Topology
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