【发布时间】:2016-03-07 02:13:42
【问题描述】:
我想创建一个包含 2 个具有 2 个不同主题的 kafkaSpout 的拓扑,并基于螺栓中的 sourceComponent 将这 2 个 spout 合并为一个流。
public class Topology {
private static final String topic1 = Real2
private static final String topic2 = Real1
public static void main(String[] args) throws AlreadyAliveException,
InvalidTopologyException, IOException {
BasicConfigurator.configure();
String zookeeper_root = "";
SpoutConfig kafkaConfig1 = new SpoutConfig(localhost:2181, topic1,
zookeeper_root, "Real1KafkaSpout");
SpoutConfig kafkaConfig2 = new SpoutConfig(localhost:2181, topic2,
zookeeper_root, "Real2KafkaSpout");
kafkaConfigRealTime.scheme = new SchemeAsMultiScheme(new StringScheme());
kafkaConfigRealTime.forceFromStart = true;
kafkaConfigHistorical.scheme = new SchemeAsMultiScheme(
new StringScheme());
kafkaConfigHistorical.forceFromStart = true;
TopologyBuilder builder = new TopologyBuilder();
builder.setSpout("Real1", new KafkaSpout(
kafkaConfig1), 2);
builder.setSpout("Real2", new KafkaSpout(
kafkaConfig2), 2);
builder.setBolt("StreamMerging", new StreamMergingBolt(), 2)
.setNumTasks(2).shuffleGrouping("Real1")
.shuffleGrouping("Real2");
Config config = new Config();
config.put("hdfs.config", yamlConf);
config.setDebug(false);
config.setMaxSpoutPending(10000);
if (args.length == 0) {
LocalCluster cluster = new LocalCluster();
cluster.submitTopology("Topology", config,
builder.createTopology());
cluster.killTopology("Topology");
cluster.shutdown();
} else {
StormSubmitter.submitTopology(args[0], config,
builder.createTopology());
}
try {
Thread.sleep(6000);
} catch (InterruptedException ex) {
ex.printStackTrace();
}
}
}
在我正在做的螺栓执行方法中
public void execute(Tuple input, BasicOutputCollector collector) {
String id = input.getSourceComponent();
System.out.println("Stream Id in StreamMergingBolt is " + "---->" + id);
}
所以我想将来自每个流的元组存储到单独的文件中 那就是我想将 Real1KafkaSpout 的元组存储到 file1 和 Real2KafkaSpout 到 file2 。我怎么能做到这一点我被打动了
【问题讨论】:
标签: java apache-kafka apache-storm