【发布时间】:2020-02-27 06:45:13
【问题描述】:
假设我有如下几列:
EMP_ID, EMP_NAME, EMP_CONTACT
1, SIDDHESH, 544949461
现在我想验证数据是否与列名架构同步。对于EMP_NAME,该列中的数据应仅为string。我在引用this 链接后尝试了下面的代码,但它在我的代码的最后一行显示错误。
package com.sample
import org.apache.spark.SparkContext
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.Row
class sample1 {
val spark = SparkSession.builder().master("local[*]").getOrCreate()
val data = spark.read.format("csv").option("header", "true").load("C:/Users/siddheshk2/Desktop/words.txt")
val originalSchema = data.schema
def validateColumns(row: Row): Row = {
val emp_id = row.getAs[String]("EMP_ID")
val emp_name = row.getAs[String]("EMP_NAME")
val emp_contact = row.getAs[String]("EMP_CONTACT")
// do checking here and populate (err_col,err_val,err_desc) with values if applicable
Row.merge(row)
}
val validateDF = data.map { row => validateColumns(row) }
}
所以,它不接受我的代码val validateDF = data.map { row => validateColumns(row) } 的最后一行。我该如何解决这个问题?或者有没有其他有效的方法可以解决我的问题?
我输入了一条无效记录(第三条),如下所示:
EMP_ID,EMP_NAME,EMP_CONTACT
1,SIDDHESH,99009809
2,asdadsa, null
sidh,sidjh,1232
在这种情况下,我已经为 id 列输入了一个 string 值(应该是一个数字),因此在检查列架构及其数据后,它应该会抛出一个错误,指出记录不匹配根据列模式。
【问题讨论】:
-
在读取数据时,您始终可以添加架构或从文件中推断架构。请查收:stackoverflow.com/questions/39926411/…
标签: scala apache-spark