【发布时间】:2016-11-29 20:41:14
【问题描述】:
我正在开发一个需要从 Kafka 读取数据的 spark 应用程序。我创建了一个 Kafka 主题,生产者在其中发布消息。我从控制台消费者验证消息已成功发布。
我编写了一个简短的 spark 应用程序来从 Kafka 读取数据,但它没有获取任何数据。 以下是我使用的代码:
def main(args: Array[String]): Unit = {
val Array(zkQuorum, group, topics, numThreads) = args
val sparkConf = new SparkConf().setAppName("SparkConsumer").setMaster("local[2]")
val ssc = new StreamingContext(sparkConf, Seconds(2))
val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
val lines = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)
process(lines) // prints the number of records in Kafka topic
ssc.start()
ssc.awaitTermination()
}
private def process(lines: DStream[String]) {
val z = lines.count()
println("count of lines is "+z)
//edit
lines.foreachRDD(rdd => rdd.map(println)
// <-- Why does this **not** print?
)
关于如何解决此问题的任何建议?
******编辑****
我用过
lines.foreachRDD(rdd => rdd.map(println)
在实际代码中也是如此,但这也不起作用。我设置了帖子中提到的保留期:Kafka spark directStream can not get data。但问题依然存在。
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming