【问题标题】:Dynamically framing column names to pass it for select while creating Spark Data frame在创建 Spark Data 框架时动态构建列名称以将其传递给 select
【发布时间】:2019-06-27 17:24:18
【问题描述】:

在过去的 2 个月里,我一直在使用 Spark 和 Scala,而且我是这项技术的新手。我将选择列(使用 regexp_replace)设置为 List [String] () 并传递给 Spark Data 框架创建,并将其抛出错误为“无法解析”。请在下面找到步骤,我已经按照并尝试过。

定义值:

Defining the column which I would like to identify in the src data frame
val col_name = "region_id"
Defining the column which will be used to replace the src data frame column from ref data frame
val surr_key_col_name = "surrogate_key"

我创建了两个数据框,如下所示

src_df

region id | region_name | region_code  
10001189  | Spain       |    SP09 8545  
10001765  | Africa      |    AF97 6754  

ref_df

region id | surrogate_key  
1189      |    2345  
1765      |    8978  

val src_df = spark.read.format("csv").option("header", "true").option("inferSchema", "true").load("s3://bucket/src_details.csv")
val ref_df = spark.read.format("csv").option("header", "true").option("inferSchema", "true").load("s3://bucket/ref_details.csv")

我正在迭代以识别我需要使用 reg match 并替换为另一个 Data Frame 列值并将其存储在 List 中以将其传递给 Data Frames select 的列

val src_header_rec = src_df.columns.toList

//Loop through source file header to identify the region_id and replace it with surrogate_id by doing a pattern match( I don't want to replace the 
for (src_header_cols <- src_header_rec) {
if (col_name == src_header_cols) {
src_column_names :+="regexp_replace("+"$"+s""""src.$src_header_cols""""+","+"$"+s""""ref.$src_header_cols""""+","+"$"+s""""ref.$surr_key_col_name""""+")"+".as("+s""""$src_header_cols""""+")"
}
else {
src_column_names :+= "src."+src_header_cols
}
}

使用上面的 for 循环在 List [String] () 中构建选择列后,我将其传递给选择列以创建 final_df

val final_df = src_df.alias("src").join(ref_df.alias("ref"), src_df(col_name)=== ref_df(col_name),"left_outer").select(src_column_names.head,src_column_names.tail:_*)

如果我直接传递列而不在数据框的选择中使用 List [String] () 我的 regexp_replace 替换工作

val final_df = src_df.alias("src").join(ref_df.alias("ref"), src_df(col_name)=== ref_df(col_name),"left_outer").select(regexp_replace($"src.region_id",$"ref.region_id",$"ref.surrogate_key").as("region_id"))

我不确定为什么当我将它作为 List [String] () 传递时它不起作用

当我在 for 循环中删除 regexp_replace 替换并将其作为数据框的 List [String] () 传递时,它可以正常工作,如下所示:

此代码与数据框选择配合得很好:

for (src_header_cols <- src_header_rec) {
if (col_name == src_header_cols) {
src_column_names :+= "ref."+surr_key_col_name
}
else {
src_column_names :+= "src."+src_header_cols
}
}

val final_df = src_df.alias("src").join(ref_df.alias("ref"), src_df(col_name)===ref_df(col_name),"left_outer").select(src_column_names.head,src_column_names.tail:_*)

我试图导出的结果/输出数据框是

final_df

region id    | region_name | region_code  
1000**2345** | Spain       |  SP09 8545  
1000**8978** | Africa      |  AF97 6754  

因此,当我尝试在 for 循环中使用 regexp_replace 作为列表构建 Spark 数据框并使用它时,它会抛出“无法解析”错误。

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    我尝试创建数据框的临时视图,并在临时视图的 select 语句中使用了相同的正则表达式。有效。请在下面找到我尝试过的代码。

    //This for loop will scan through my header list and whichever column matches it frames regexp for those columns.So, the region_id from the Data Frame header matches the variable value that I have defined.
    for (src_header_cols <- src_header_rec) {
        if (col_name == src_header_cols) {
    
            src_column_names :+= "regexp_replace(src."+s"$src_header_cols"+",ref."+s"$ref_col_name"+",ref."+s"$surr_key_col_name"+")"+s" $src_header_cols"
        }
        else {
            src_column_names :+= "src."+src_header_cols
        }
    } 
    
    //Creating Temporary view to apply SQL queries on it
    src_df.createOrReplaceTempView("src")
    ref_df.createOrReplaceTempView("ref")
    
    //Framing SQL statements to be passed while querying
    val selectExpr_1 = "select "+src_column_names.mkString(",")
    val selectExpr_2 = selectExpr_1+" from src left outer join ref on(src."+s"$col_name"+" = ref."+s"$ref_col_name"+")"
    
    // Creating a final Data Frame using the SQL statement created
    val src_policy_masked_df = spark.sql(s"$selectExpr_2")
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-06-11
      • 1970-01-01
      • 2017-07-03
      • 2022-11-22
      • 1970-01-01
      • 1970-01-01
      • 2017-12-10
      • 1970-01-01
      相关资源
      最近更新 更多