【发布时间】:2016-05-30 21:09:37
【问题描述】:
我正在阅读来自 2 个不同主题的 2 个不同文件的 kafka 行。行示例:
例如:
文件 1:2015-04-15T18:44:14+01:00,192.168.11.42,%ASA-2-106007:
文件2:"04/15/2012","18:44:14",,"Start","Unknown","Unknown",,"192.168.63.128","444","2","7","192.168.63.128",,,,,,,,,,,,,,,,,
我可以从两个不同主题的火花中阅读。代码如下:
SparkConf sparkConfig = new SparkConf().setAppName("KafkaStreaming").setMaster("local[5]");
JavaStreamingContext jsc = new JavaStreamingContext(sparkConfig,Durations.seconds(5));
final HiveContext sqlContext = new HiveContext(jsc.sc());
JavaPairReceiverInputDStream<String, String> messages = KafkaUtils.createStream(jsc,
prop.getProperty("zookeeper.connect"),
prop.getProperty("group.id"),
topicMap
);
JavaDStream<String> lines = messages.map(new Function<Tuple2<String, String>, String>() {
private static final long serialVersionUID = 1L;
public String call(Tuple2<String, String> tuple2) {
return tuple2._2();
}
});
我现在看到的问题是:
lines rdd 包含两个明显的行。我如何分离或找出哪些记录来自哪个主题或哪个文件。
这背后的原因是我想要为即将到来的不同主题应用不同的逻辑。但是 rdd 有时间的所有行
感谢任何建议
【问题讨论】:
标签: apache-spark apache-kafka kafka-consumer-api