【发布时间】:2018-06-17 18:30:32
【问题描述】:
嗯,我是 spark 和 scala 的新手,一直在尝试在 spark 中实现数据清理。下面的代码检查一列的缺失值并将其存储在 outputrdd 中并运行循环以计算缺失值。当文件中只有一个缺失值时,代码运行良好。由于 hdfs 不允许在同一位置再次写入,因此如果有多个缺失值,它会失败。完成计算所有事件的缺失值后,您能否协助将 finalrdd 写入特定位置。
def main(args: Array[String]) {
val conf = new SparkConf().setAppName("app").setMaster("local")
val sc = new SparkContext(conf)
val sqlContext = new org.apache.spark.sql.SQLContext(sc)
val files = sc.wholeTextFiles("/input/raw_files/")
val file = files.map { case (filename, content) => filename }
file.collect.foreach(filename => {
cleaningData(filename)
})
def cleaningData(file: String) = {
//headers has column headers of the files
var hdr = headers.toString()
var vl = hdr.split("\t")
sqlContext.clearCache()
if (hdr.contains("COLUMN_HEADER")) {
//Checks for missing values in dataframe and stores missing values' in outputrdd
if (!outputrdd.isEmpty()) {
logger.info("value is zero then performing further operation")
val outputdatetimedf = sqlContext.sql("select date,'/t',time from cpc where kwh = 0")
val outputdatetimerdd = outputdatetimedf.rdd
val strings = outputdatetimerdd.map(row => row.mkString).collect()
for (i <- strings) {
if (Coddition check) {
//Calculates missing value and stores in finalrdd
finalrdd.map { x => x.mkString("\t") }.saveAsTextFile("/output")
logger.info("file is written in file")
}
}
}
}
}
}``
【问题讨论】:
标签: scala apache-spark spark-dataframe rdd