【发布时间】:2016-10-12 10:40:05
【问题描述】:
如何将Spark Streaming 中的RDD 转换为DataFrame,而不仅仅是Spark?
我看到了这个例子,但它需要SparkContext。
val sqlContext = new SQLContext(sc)
import sqlContext.implicits._
rdd.toDF()
就我而言,我有StreamingContext。然后我应该在foreach 中创建SparkContext 吗?看起来太疯狂了……那么,如何处理这个问题呢?我的最终目标(如果可能有用的话)是使用rdd.toDF.write.format("json").saveAsTextFile("s3://iiiii/ttttt.json"); 将DataFrame 保存在Amazon S3 中,如果不将RDD 转换为DataFrame(据我所知),这是不可能的。
myDstream.foreachRDD { rdd =>
val conf = new SparkConf().setMaster("local").setAppName("My App")
val sc = new SparkContext(conf)
val sqlContext = new SQLContext(sc)
import sqlContext.implicits._
rdd.toDF()
}
【问题讨论】:
-
@Shankar:他在哪里定义 AWS 访问密钥?
-
foreachRDD中写入的任何内容都会在驱动程序中执行,因此您可以创建sqlContext并将rdd转换为DF,然后写入S3。 -
@Shankar:我仍然误解:我应该在
foreachRDD之外创建 StreamingContext 和 SparkContext 吗?在您发布的示例中,我找不到sqlContext的定义位置。我尝试重现此示例,但它给了我一个错误,即找不到sqlContext。我不想让事情变得过于复杂,这就是为什么我会询问最简单的解决方案。
标签: scala apache-spark spark-streaming rdd