【问题标题】:How to validate large csv file either column wise or row wise in spark dataframe如何在 Spark 数据框中逐列或逐行验证大型 csv 文件
【发布时间】:2020-08-23 14:15:45
【问题描述】:
我有一个 10GB 或更大的大型数据文件,包含 150 列,我们需要在其中使用不同的规则验证其每个数据(数据类型/格式/空值/域值/主键..),最后创建 2 个输出文件一个是成功数据,另一个是带有错误详细信息的错误数据。如果任何列在第一次出现错误,我们需要移动错误文件中的行,不需要进一步验证。
我正在读取 spark 数据帧中的文件,我们是按列还是按行验证它,我们通过哪种方式获得了最佳性能?
【问题讨论】:
标签:
apache-spark
apache-spark-sql
【解决方案1】:
回答你的问题
我正在读取 spark 数据帧中的文件,我们是按列还是按行验证它,我们通过哪种方式获得了最佳性能?
DataFrame 是一个分布式数据集合,它被组织为分布在集群中的一组行,并且在 spark 中定义的大部分转换都应用于在 Row object 上工作的行。
Psuedo code
import spark.implicits._
val schema = spark.read.csv(ip).schema
spark.read.textFile(inputFile).map(row => {
val errorInfo : Seq[(Row,String,Boolean)] = Seq()
val data = schema.foreach(f => {
// f.dataType //get field type and have custom logic on field type
// f.name // get field name i.e., column name
// val fieldValue = row.getAs(f.name) //get field value and have check's on field value on field type
// if any error in field value validation then populate @errorInfo info object i.e (row,"error_info",false)
// otherwise i.e (row,"",true)
})
data.filter(x => x._3).write.save(correctLoc)
data.filter(x => !x._3).write.save(errorLoc)
})