代码之家  ›  专栏  ›  技术社区  ›  1pluszara

使用spark数据帧进行现场数据验证

  •  1
  • 1pluszara  · 技术社区  · 9 年前

    我需要检查列是否有错误,并且必须生成两个输出文件。 我正在使用Apache Spark 2.0,我希望以一种高效的方式来实现这一点。

    Schema Details
    ---------------
    EMPID - (NUMBER)
    ENAME - (STRING,SIZE(50))
    GENDER - (STRING,SIZE(1))
    
    Data
    ----
    EMPID,ENAME,GENDER
    1001,RIO,M
    1010,RICK,MM
    1015,123MYA,F
    

    我的预期输出文件应如下所示:

    1.
    EMPID,ENAME,GENDER
    1001,RIO,M
    1010,RICK,NULL
    1015,NULL,F
    
    2.
    EMPID,ERROR_COLUMN,ERROR_VALUE,ERROR_DESCRIPTION
    1010,GENDER,"MM","OVERSIZED"
    1010,GENDER,"MM","VALUE INVALID FOR GENDER"
    1015,ENAME,"123MYA","NAME SHOULD BE A STRING"
    

    1 回复  |  直到 9 年前
        1
  •  8
  •   Nilanjan Sarkar    9 年前

    我还没有真正使用Spark 2.0,所以我将尝试用Spark 1.6中的解决方案来回答您的问题。

     // Load you base data
    val input  = <<you input dataframe>>
    
    //Extract the schema of your base data
    val originalSchema = input.schema
    
    // Modify you existing schema with you additional metadata fields
    val modifiedSchema= originalSchema.add("ERROR_COLUMN", StringType, true)
                                      .add("ERROR_VALUE", StringType, true)
                                      .add("ERROR_DESCRIPTION", StringType, true)
    
    // write a custom validation function                                 
    def validateColumns(row: Row): Row = {
    
    var err_col: String = null
    var err_val: String = null
    var err_desc: String = null
    val empId = row.getAs[String]("EMPID")
    val ename = row.getAs[String]("ENAME")
    val gender = row.getAs[String]("GENDER")
    
    // do checking here and populate (err_col,err_val,err_desc) with values if applicable
    
    Row.merge(row, Row(err_col),Row(err_val),Row(err_desc))
    }
    
    // Call you custom validation function
    val validateDF = input.map { row => validateColumns(row) }  
    
    // Reconstruct the DataFrame with additional columns                      
    val checkedDf = sqlContext.createDataFrame(validateDF, newSchema)
    
    // Filter out row having errors
    val errorDf = checkedDf.filter($"ERROR_COLUMN".isNotNull && $"ERROR_VALUE".isNotNull && $"ERROR_DESCRIPTION".isNotNull)
    
    // Filter our row having no errors
    val errorFreeDf = checkedDf.filter($"ERROR_COLUMN".isNull && !$"ERROR_VALUE".isNull && !$"ERROR_DESCRIPTION".isNull)
    

    推荐文章