【发布时间】:2016-09-09 13:09:15
【问题描述】:
您好,我正在尝试使用 Apache Spark Streaming 从 Twitter 读取推文并尝试转换为 DataFrame。我有我在下面粘贴的方法。但是,我无法获得正确的方法。一些指针将受到欢迎。
正如您所见,在 foreach 中转换为 DF 并没有从 tweetStream 中获得单个 DF。我可能有错误的方法,因为我是新手。我该如何处理?
val tweetStream = TwitterUtils.createStream(ssc, Utils.getAuth).filter(status=>status.getLang=="en")
.map(status=>gson.toJson(status))
val sqlContext = new org.apache.spark.sql.SQLContext(sc)
import sqlContext.implicits._
tweetStream.foreachRDD({status=>val DF = status.toDF()})
【问题讨论】:
-
我正在考虑在循环中使用 DF.merge() 来获取在 foreachRDD 中计算的整个 DF{}
标签: scala apache-spark bigdata