【问题标题】:Spark remove special characters from column name read from a parquet file [duplicate]Spark从镶木地板文件中读取的列名中删除特殊字符[重复]
【发布时间】:2021-05-06 18:10:59
【问题描述】:

我有使用以下 spark 命令读取的镶木地板文件

lazy val out = spark.read.parquet("/tmp/oip/logprint_poc/feb28eb24ffe44cab60f2832a98795b1.parquet")

很多列的列名都有特殊字符“(”。比如WA_0_DWHRPD_Purge_Date_(TOD)WA_0_DWHRRT_Record_Type_(80=Index)我怎样才能去掉这个特殊字符。

我的最终目标是删除这些特殊字符并使用以下命令写回镶木地板文件

df_hive.write.format("parquet").save("hdfs:///tmp/oip/logprint_poc_cleaned/")

另外,我正在使用 Scala spark shell。 我是新手,我看到了类似的问题,但在我的情况下没有任何效果。任何帮助表示赞赏。

【问题讨论】:

标签: apache-spark parquet


【解决方案1】:

您可以做的第一件事就是将 parquet 文件读入数据框。

val out = spark.read.parquet("/tmp/oip/logprint_poc/feb28eb24ffe44cab60f2832a98795b1.parquet")

创建数据框后,尝试获取数据框的架构并对其进行解析以删除所有特殊字符,如下所示:

import org.apache.spark.sql.functions._
val schema = StructType(out.schema.map(
          x => StructField(x.name.toLowerCase().replace(" ", "_").replace("#", "").replace("-", "_").replace(")", "").replace("(", "").trim(),
            x.dataType, x.nullable)))

现在您可以通过指定您创建的架构从 parquet 文件中读取数据。

val newDF = spark.read.format("parquet").schema(schema).load("/tmp/oip/logprint_poc/feb28eb24ffe44cab60f2832a98795b1.parquet")

现在您可以使用已清理的列名称继续保存数据框。

df_hive.write.format("parquet").save("hdfs:///tmp/oip/logprint_poc_cleaned/")

【讨论】:

    猜你喜欢
    • 2017-01-22
    • 1970-01-01
    • 1970-01-01
    • 2018-06-02
    • 2020-08-15
    • 1970-01-01
    • 2021-08-27
    • 2021-01-12
    • 1970-01-01
    相关资源
    最近更新 更多