【问题标题】:Scala spark: how to use dataset for a case class with the schema has snake_case?Scala spark:如何将数据集用于具有snake_case的模式的案例类?
【发布时间】:2018-09-25 22:40:38
【问题描述】:

我有以下案例类:

case class User(userId: String)

以及以下架构:

+--------------------+------------------+
|            col_name|         data_type|
+--------------------+------------------+
|             user_id|            string|
+--------------------+------------------+

当我尝试使用spark.read.table("MyTable").as[User]DataFrame 转换为键入的Dataset[User] 时,出现字段名称不匹配的错误:

Exception in thread "main" org.apache.spark.sql.AnalysisException:
    cannot resolve ''`user_id`' given input columns: [userId];;

是否有任何简单 的方法来解决这个问题而不会破坏 scala 成语并将我的字段命名为 user_id?当然,我的真实表的字段多得多,而且我的案例类/表也多得多,所以为每个案例类手动定义一个Encoder是不可行的(而且我对宏的了解不够,所以这是不可能的;尽管如果存在的话,我很乐意使用它!)。

我觉得我缺少一个非常明显的“将蛇形大小写转换为camelCase = true”选项,因为它几乎存在于我使用过的任何ORM中。

【问题讨论】:

  • 我觉得我错过了一个非常明显的“将 snake_case 转换为 camelCase=true” - 你没有。如果我没记错的话,有一些旧的 JIRA 票针对类似的东西,但现在,你必须重命名。
  • @user6910411 Bummer :( 如果你用 JIRA 票回答,我会接受答案。
  • @Gal 三年后,你找到更好的解决方案了吗?
  • @DanielR 不幸的是,没有。如果我的 case class 字段代表一个 spark 表,我就辞职了。

标签: scala apache-spark apache-spark-dataset


【解决方案1】:
scala> val df = Seq(("Eric" ,"Theodore", "Cartman"), ("Butters", "Leopold", "Stotch")).toDF.select(concat($"_1", lit(" "), ($"_2")) as "first_and_middle_name", $"_3" as "last_name")
df: org.apache.spark.sql.DataFrame = [first_and_middle_name: string, last_name: string]

scala> df.show
+---------------------+---------+
|first_and_middle_name|last_name|
+---------------------+---------+
|        Eric Theodore|  Cartman|
|      Butters Leopold|   Stotch|
+---------------------+---------+


scala> val ccnames = df.columns.map(sc => {val ccn = sc.split("_")
    | (ccn.head +: ccn.tail.map(_.capitalize)).mkString
    | })
ccnames: Array[String] = Array(firstAndMiddleName, lastName)

scala> df.toDF(ccnames: _*).show
+------------------+--------+
|firstAndMiddleName|lastName|
+------------------+--------+
|     Eric Theodore| Cartman|
|   Butters Leopold|  Stotch|
+------------------+--------+

编辑:这会有帮助吗?定义一个带 loader 的函数:String => DataFrame 和 path:String。

scala> val parquetloader = spark.read.parquet _
parquetloader: String => org.apache.spark.sql.DataFrame = <function1>

scala> val tableloader = spark.read.table _
tableloader: String => org.apache.spark.sql.DataFrame = <function1>

scala> val textloader = spark.read.text _
textloader: String => org.apache.spark.sql.DataFrame = <function1>

// csv loader and others

def snakeCaseToCamelCaseDataFrameColumns(path: String, loader: String => DataFrame): DataFrame = {
  val ccnames = loader(path).columns.map(sc => {val ccn = sc.split("_")
    (ccn.head +: ccn.tail.map(_.capitalize)).mkString
    })
  df.toDF(ccnames: _*)
}

scala> :paste
// Entering paste mode (ctrl-D to finish)

def snakeCaseToCamelCaseDataFrameColumns(path: String, loader: String => DataFrame): DataFrame = {
      val ccnames = loader(path).columns.map(sc => {val ccn = sc.split("_")
        (ccn.head +: ccn.tail.map(_.capitalize)).mkString
        })
      df.toDF(ccnames: _*)
    }

// Exiting paste mode, now interpreting.

snakeCaseToCamelCaseDataFrameColumns: (path: String, loader: String => org.apache.spark.sql.DataFrame)org.apache.spark.sql.DataFrame

val oneDF = snakeCaseToCamelCaseDataFrameColumns(tableloader("/path/to/table"))
val twoDF = snakeCaseToCamelCaseDataFrameColumns(parquetloader("/path/to/parquet/file"))

【讨论】:

  • 正如我所说,“当然,我的真实表有很多更多的字段,并且我有更多的案例类/表,因此手动操作是不可行的”。
猜你喜欢
  • 2021-12-30
  • 1970-01-01
  • 1970-01-01
  • 2017-04-25
  • 1970-01-01
  • 2017-04-11
  • 1970-01-01
  • 1970-01-01
  • 2017-10-14
相关资源
最近更新 更多