【问题标题】:How to convert kafka stream to spark RDD or Spark Dataframe如何将 kafka 流转换为 spark RDD 或 Spark Dataframe
【发布时间】:2016-05-12 16:50:00
【问题描述】:

我尝试从 Kafka 加载数据成功,但无法转换为 spark RDD,

val kafkaParams = Map("metadata.broker.list" -> "IP:6667,IP:6667")
val offsetRanges = Array(
    OffsetRange("first_topic", 0,1,1000)
  )
val ssc = new StreamingContext(new SparkConf, Seconds(60))
val stream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc, kafkaParams, topics)

现在我怎样才能读取这个流对象???我的意思是将其转换为 Spark Dataframe 并执行一些计算

我尝试转换为数据框

    stream.foreachRDD { rdd =>
     println("Hello")
      import sqlContext.implicits._
      val dataFrame = rdd.map {case (key, value) => Row(key, value)}.toDf()
    }

但是 toDf 不起作用错误:值 toDf 不是 org.apache.spark.rdd.RDD[org.apache.spark.sql.Row] 的成员

【问题讨论】:

    标签: scala apache-spark apache-kafka spark-streaming


    【解决方案1】:

    它很旧,但我认为您在从行创建 df 时忘记添加架构:

    val df =  sc.parallelize(List(1,2,3)).toDF("a")
    val someRDD = df.rdd
    val newDF = spark.createDataFrame(someRDD, df.schema)
    

    (在 spark-shell 2.2.0 中测试)

    【讨论】:

      【解决方案2】:
      val kafkaParams = Map("metadata.broker.list" -> "IP:6667,IP:6667")
      val offsetRanges = Array(
          OffsetRange("first_topic", 0,1,1000)
        )
      val ssc = new StreamingContext(new SparkConf, Seconds(60))
      val stream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc, kafkaParams, topics) 
      
      val lines = stream.map(_.value)
      val words = lines.flatMap(_.split(" ")).print()   //def createDataFrame(words: RDD[Row], Schema: StructType)
      
      // Start your computation then
      ssc.start()
      ssc.awaitTermination()
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-06-28
        • 2018-03-05
        • 2017-06-02
        • 1970-01-01
        • 2020-02-28
        • 2017-06-13
        • 2018-10-21
        • 1970-01-01
        相关资源
        最近更新 更多