【发布时间】: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