【发布时间】:2017-01-25 23:52:14
【问题描述】:
我是 Spark 流媒体 的新手,我正在实施小练习,例如从 kafka 发送 XML 数据,并且需要接收 >通过火花流传输数据。我尝试了所有可能的方式..但每次我得到空值。
Kafka端没有问题,唯一的问题是从Spark端接收流数据。
这是我如何实现的代码:
package com.package;
import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.streaming.Duration;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
public class SparkStringConsumer {
public static void main(String[] args) {
SparkConf conf = new SparkConf()
.setAppName("kafka-sandbox")
.setMaster("local[*]");
JavaSparkContext sc = new JavaSparkContext(conf);
JavaStreamingContext ssc = new JavaStreamingContext(sc, new Duration(2000));
Map<String, String> kafkaParams = new HashMap<>();
kafkaParams.put("metadata.broker.list", "localhost:9092");
Set<String> topics = Collections.singleton("mytopic");
JavaPairInputDStream<String, String> directKafkaStream = KafkaUtils.createDirectStream(ssc,
String.class, String.class, StringDecoder.class, StringDecoder.class, kafkaParams, topics);
directKafkaStream.foreachRDD(rdd -> {
System.out.println("--- New RDD with " + rdd.partitions().size()
+ " partitions and " + rdd.count() + " records");
rdd.foreach(record -> System.out.println(record._2));
});
ssc.start();
ssc.awaitTermination();
}
}
我正在使用以下版本:
**动物园管理员 3.4.6
Scala 2.11
火花 2.0
卡夫卡 0.8.2**
【问题讨论】:
标签: hadoop apache-spark streaming apache-kafka spark-streaming