【问题标题】:Read nested columns from CSV file and assign a schema to the dataframe从 CSV 文件中读取嵌套列并将架构分配给数据框
【发布时间】:2021-05-19 10:00:41
【问题描述】:

我尝试读取包含嵌套列的 csv 文件。

例子:

name,age,addresses_values
person_1,30,["France","street name",75000]

阅读时我厌倦了分配如下模式:

csv_schema = StructType([
            StructField('name', StringType(), True),
            StructField('age', LongType(), True),
            StructField('addresses_values', StructType([
                    StructField('country', StringType(), True),
                    StructField('street', StringType(), True),
                   StructField('ZipCode', StringType(), True),
                   ]), True),
        ])

path = "file:///path_to_my_file"

dataset_df = spark.read.csv(path=path, header=True,schema=csv_schema)

引发此异常:

pyspark.sql.utils.AnalysisException:CSV 数据源不支持 structcountry:string,street:string,ZipCode:string 数据类型。;

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql


    【解决方案1】:

    有点棘手,但这里有一种解析这些 CSV 值的方法。您需要指定quote="[" 才能读取addresses_values 列,因为它包含逗号。然后使用from_json 将其解析为字符串数组,最后从数组元素创建所需的结构:

    from pyspark.sql import functions as F
    
    df = spark.read.csv(input_path, header=True, quote="[")
    
    df1 = df.withColumn(
        "addresses_values",
        F.from_json(
            F.concat(F.lit("["), F.col("addresses_values")),
            "array<string>"
        )
    ).withColumn(
        "addresses_values",
        F.struct(
            F.col("addresses_values")[0].alias("country"),
            F.col("addresses_values")[1].alias("street"),
            F.col("addresses_values")[2].alias("ZipCode"),
        )
    )
    
    df1.show(truncate=False)
    
    #+--------+---+----------------------------+
    #|name    |age|addresses_values            |
    #+--------+---+----------------------------+
    #|person_1|30 |[France, street name, 75000]|
    #+--------+---+----------------------------+
    
    df1.printSchema()
    
    #root
    # |-- name: string (nullable = true)
    # |-- age: string (nullable = true)
    # |-- addresses_values: struct (nullable = false)
    # |    |-- country: string (nullable = true)
    # |    |-- street: string (nullable = true)
    # |    |-- ZipCode: string (nullable = true)
    

    【讨论】:

      【解决方案2】:

      这是一种令人讨厌的解析方式。灵感来自 thisthis

      import pyspark.sql.functions as F
      
      df = spark.read.text('arr.csv') \
          .filter("value != 'name,age,addresses_values'") \
          .select(F.split('value', ',(?=(?:[^\[\]]*\[[^\[\]]*\])*[^\[\]]*$)').alias('value')) \
          .selectExpr('value[0] name', 'value[1] age', "split(trim('[]' from value[2]), ',(?=(?:[^\"]*\"[^\"]*\")*[^\"]*$)') addresses_values") \
          .selectExpr('name', 'age', 
              """struct(trim('"' from addresses_values[0]) as country,
                        trim('"' from addresses_values[1]) as street,
                        addresses_values[2] as zipcode)
                 as addresses_values
              """)
      
      df.show(truncate=False)
      +--------+---+----------------------------+
      |name    |age|addresses_values            |
      +--------+---+----------------------------+
      |person_1|30 |[France, street name, 75000]|
      +--------+---+----------------------------+
      
      df.printSchema()
      root
       |-- name: string (nullable = true)
       |-- age: string (nullable = true)
       |-- addresses_values: struct (nullable = false)
       |    |-- country: string (nullable = true)
       |    |-- street: string (nullable = true)
       |    |-- zipcode: string (nullable = true)
      

      【讨论】:

        猜你喜欢
        • 2013-08-07
        • 1970-01-01
        • 2015-06-23
        • 1970-01-01
        • 2021-02-22
        • 2017-02-16
        • 2021-08-14
        • 2021-02-06
        • 1970-01-01
        相关资源
        最近更新 更多