【问题标题】:Spark 2.0.x dump a csv file from a dataframe containing one array of type stringSpark 2.0.x 从包含一个字符串类型数组的数据帧中转储 csv 文件
【发布时间】:2016-11-04 15:15:09
【问题描述】:

我有一个数据框df,其中包含一列数组类型

df.show() 看起来像

|ID|ArrayOfString|Age|Gender|
+--+-------------+---+------+
|1 | [A,B,D]     |22 | F    |
|2 | [A,Y]       |42 | M    |
|3 | [X]         |60 | F    |
+--+-------------+---+------+

我尝试将 df 转储到 csv 文件中,如下所示:

val dumpCSV = df.write.csv(path="/home/me/saveDF")

由于ArrayOfString 列,它无法正常工作。我得到了错误:

CSV数据源不支持数组字符串数据类型

如果我删除 ArrayOfString 列,该代码将起作用。但我需要保留ArrayOfString

转储包含 ArrayOfString 列的 csv 数据帧的最佳方法是什么(ArrayOfString 应作为 CSV 文件中的一列转储)

【问题讨论】:

    标签: arrays csv apache-spark


    【解决方案1】:

    如果您已经知道哪些字段包含数组,则不需要 UDF。你可以简单地使用 Spark 的 cast 函数:

    import org.apache.spark.sql.functions._
    val dumpCSV = df.withColumn("ArrayOfString", col("ArrayOfString").cast("string"))
                    .write
                    .csv(path="/home/me/saveDF")
    

    希望对您有所帮助。

    【讨论】:

    • 这应该是最佳答案,它使用内置函数来解决问题。
    【解决方案2】:

    您收到此错误的原因是 csv 文件格式不支持数组类型,您需要将其表示为字符串才能保存。

    尝试以下方法:

    import org.apache.spark.sql.functions._
    
    val stringify = udf((vs: Seq[String]) => vs match {
      case null => null
      case _    => s"""[${vs.mkString(",")}]"""
    })
    
    df.withColumn("ArrayOfString", stringify($"ArrayOfString")).write.csv(...)
    

    import org.apache.spark.sql.Column
    
    def stringify(c: Column) = concat(lit("["), concat_ws(",", c), lit("]"))
    
    df.withColumn("ArrayOfString", stringify($"ArrayOfString")).write.csv(...)
    

    【讨论】:

    • 您好,非常感谢您的回答。我明白这些线的作用。但是我对语法 s"""[${vs.mkString(",")}]""" 有点困惑,你能解释一下关于 s 和三元组 """ 的更多信息吗?谢谢。跨度>
    • 嗯,感谢您发给我的文档,我更好地理解了“s”的含义。但是我仍然不明白为什么 3 引号。为什么我不能写 s"[${vs.mkString(",")}]" 顺便说一句,使用 1 引号也适用于我。那么为什么是 3 个引号呢?
    【解决方案3】:

    Pyspark 实现。

    在此示例中,在保存之前将字段 column_as_array 更改为 column_as_string

    from pyspark.sql.functions import udf
    from pyspark.sql.types import StringType
    
    def array_to_string(my_list):
        return '[' + ','.join([str(elem) for elem in my_list]) + ']'
    
    array_to_string_udf = udf(array_to_string, StringType())
    
    df = df.withColumn('column_as_str', array_to_string_udf(df["column_as_array"]))
    

    然后你可以在保存之前删除旧列(数组类型)。

    df.drop("column_as_array").write.csv(...)
    

    【讨论】:

    • 我有两列“Antecedent”和“Consequent”,其中有列表作为输入。我怎样才能修改这段代码来做同样的事情。
    • 它对我有用(甚至可以通过创建一个列表来自动化它 all_mappings = [x for (x,y) in df_p300.dtypes if y == 'map<string,int>' ] )。小细节:第一块代码的最后一行应该是df["column_as_array"]而不是d["column_as_array"]
    • 有人可以像@LionelTrebuchon 提到的那样将d 修复为df 吗? SO 不允许编辑一个字符,但是这个例子有效,但是这个错字。
    • 字符错字已修复,感谢 LionelTrebuchon 和 @zeh
    【解决方案4】:

    这是一种将DataFrame 的所有ArrayType(任何基础类型)列转换为StringType 列的方法:

    def stringifyArrays(dataFrame: DataFrame): DataFrame = {
      val colsToStringify = dataFrame.schema.filter(p => p.dataType.typeName == "array").map(p => p.name)
      colsToStringify.foldLeft(dataFrame)((df, c) => {
        df.withColumn(c, concat(lit("["), concat_ws(", ", col(c).cast("array<string>")), lit("]")))
      })
    }
    

    此外,它不使用 UDF。

    【讨论】:

      【解决方案5】:

      CSV 不是理想的导出格式,但如果您只是想直观地检查您的数据,这将有效 [Scala]。快速而肮脏的解决方案。

      case class example ( id: String, ArrayOfString: String, Age: String, Gender: String)
      
      df.rdd.map{line => example(line(0).toString, line(1).toString, line(2).toString , line(3).toString) }.toDF.write.csv("/tmp/example.csv")
      

      【讨论】:

        【解决方案6】:

        回答 DreamerP 的问题(来自其中一位 cmets):

        from pyspark.sql.functions import udf
        from pyspark.sql.types import StringType
        
        def array_to_string(my_list):
            return '[' + ','.join([str(elem) for elem in my_list]) + ']'
        
        array_to_string_udf = udf(array_to_string, StringType())
        
        df = df.withColumn('Antecedent_as_str', array_to_string_udf(df["Antecedent"]))
        df = df.withColumn('Consequent_as_str', array_to_string_udf(df["Consequent"]))
        df = df.drop("Consequent")
        df = df.drop("Antecedent")
        df.write.csv("foldername")
        

        【讨论】:

          猜你喜欢
          • 2020-01-02
          • 1970-01-01
          • 1970-01-01
          • 2017-09-12
          • 1970-01-01
          • 1970-01-01
          • 2023-03-18
          • 2013-09-28
          • 2012-07-24
          相关资源
          最近更新 更多