【发布时间】:2017-05-04 06:01:00
【问题描述】:
我正在使用 Scala 开发 Spark 应用程序,并且想知道将其并行化并在 Hadoop 集群上运行的最佳方法。我的代码将从 HDFS 文件中读取每一行并对其进行解析并生成多个记录(对于每一行),我将这些记录存储为案例类。我已经在 getElem() 方法中编写了完整的逻辑并按预期工作。
现在,我想计算所有输入记录的逻辑并将响应存储到 HDFS 位置。
请让我知道我将如何处理 spark 并合并为输入生成的所有相应输出记录并写入 HDFS。
object testing extends Serializable {
var recordArray=Array[Record]();
def main(args:Array[String])
{
val conf = new SparkConf().setAppName("jsonParsing").setMaster("local")
val sc = new SparkContext(conf)
val sqlContext= new SQLContext(sc)
val input=sc.textFile("hdfs://loc/data.txt")
// input.collect().foreach(println)
input.map(data=>getElem(parse(data,false),sc,sqlContext))
}
//method definition
def getElem(json:JValue)={
// Parses the json and creates array of datasets for each input record and stores the data in case class
val x= Record("xxxx","xxxx","xxxx","xxxx","xxxx","xxxx","xxxx","xxxx","xxxx","xxxx")
}
case class Record(summary_key: String, key: String,array_name_position:Int,Parent_Level_1:String,Parent_level_2:String,Parent_Level_3:String,Parent_level_4:String,Parent_level_5:String,
param_name_position:Integer,Array_name:String,paramname:String,paramvalue:String)
}
【问题讨论】:
-
代码示例不完整。
recordArray从未使用过,getElem未指定返回类型(在您发布的代码中,它仅返回Unit)。它也有错误的签名,您在地图中将SparkContext和SQLContext传递给它,但在定义中它只接受JValue。parse从未被解释过。您能否发布一个工作示例来说明您的问题并准确指出它不工作的原因?
标签: scala apache-spark dataframe apache-spark-sql rdd