【发布时间】: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