【问题标题】:Why does my Spark Streaming application not print the number of records from Kafka (using count operator)?为什么我的 Spark Streaming 应用程序不打印来自 Kafka 的记录数(使用计数运算符)?
【发布时间】: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


    【解决方案1】:

    您的 process 是 DStream 管道的延续,带有 no 输出运算符,该运算符在每个批处理间隔执行管道。

    你可以通过阅读count operator的签名“看到”它:

    count(): DStream[Long]
    

    引用count's scaladoc:

    返回一个新的 DStream,其中每个 RDD 都有一个元素,该元素是通过计算此 DStream 的每个 RDD 生成的。

    因此,您有一个 Kafka 记录的 dstream,您将其转换为单个值的 dstream(是 count 的结果)。将其输出(到控制台或任何其他接收器)并不多。

    你必须使用官方文档Output Operations on DStreams中描述的输出操作符来结束管道:

    输出操作允许 DStream 的数据被推送到外部系统,如数据库或文件系统。由于输出操作实际上允许外部系统使用转换后的数据,因此它们会触发所有 DStream 转换的实际执行(类似于 RDD 的操作)。

    (低级) 输出运算符将输入 dstream 注册为输出 dstream,以便开始执行。 Spark Streaming 的 DStream 在设计上没有作为输出 dstream 的概念。 DStreamGraph 知道并能够区分输入和输出 dstream。

    【讨论】:

    • 我在实际代码中也使用了输出运算符。我编辑了原始问题以表明问题仍然存在。现在有什么建议吗?
    • 呵呵,你正在将一个无输出运算符“交易”到另一个 :) 你能先使用lines.count().print 而不是lines.count() 吗?我相信您会在控制台上打印出 10 条记录。至于RDD的情况,请使用rdd.foreach(println)(不是rdd.map(println),这是一个转换)。玩得开心! :)
    猜你喜欢
    • 1970-01-01
    • 2017-04-10
    • 2017-10-24
    • 1970-01-01
    • 2017-04-26
    • 2016-11-03
    • 1970-01-01
    • 1970-01-01
    • 2017-04-02
    相关资源
    最近更新 更多