【问题标题】:Pyspark Column Values are getting shifted automatically while creating DataFramePyspark 列值在创建 DataFrame 时自动移动
【发布时间】:2020-08-06 14:04:00
【问题描述】:

我正在尝试使用以下嵌套模式手动创建 pyspark 数据框 -

schema = StructType([
    StructField('fields', ArrayType(StructType([
        StructField('source', StringType()), 
        StructField('sourceids', ArrayType(IntegerType()))]))), 
  StructField('first_name',StringType()), 
  StructField('last_name',StringType()), 
  StructField('kare_id',StringType()),
  StructField('match_key',ArrayType(StringType()))
])

我正在使用下面的代码来创建一个使用此架构的数据框 -

row = [Row(fields=[Row(
                    source='BCONNECTED', 
                    sourceids=[10,202,30]), 
                Row(
                    source='KP', 
                    sourceids=[20,30,40])],first_name='Christopher', last_name='Nolan', kare_id='kare1', match_key=['abc','abcd']), 
        Row(fields=[
                Row(
                    source='BCONNECTED', 
                    sourceids=[20,304,5,6]), 
                Row(
                    source='KP',  
                    sourceids=[40,50,60])],first_name='Michael', last_name='Caine', kare_id='kare2', match_key=['ncnc','cncnc'])]

content = spark.createDataFrame(sc.parallelize(row), schema=schema)
content.printSchema()

Schema 打印正确,但是当我执行 content.show() 时,我可以看到 kare_id 和 last_name 列的值已交换。

+--------------------+-----------+---------+-------+-------------+
|              fields| first_name|last_name|kare_id|    match_key|
+--------------------+-----------+---------+-------+-------------+
|[[BCONNECTED, [10...|Christopher|    kare1|  Nolan|  [abc, abcd]|
|[[BCONNECTED, [20...|    Michael|    kare2|  Caine|[ncnc, cncnc]|
+--------------------+-----------+---------+-------+-------------+

【问题讨论】:

    标签: dataframe apache-spark pyspark databricks


    【解决方案1】:

    PySpark 使用字典顺序对列名上的 Row 对象进行排序。因此,数据中列的顺序将是fields, first_name, kare_id, last_name, match_key

    Spark 然后将每个列名与导致不匹配的数据相关联。修复方法是交换 last_namekare_id 的架构条目,如下所示:

    schema = StructType([
        StructField('fields', ArrayType(StructType([
        StructField('source', StringType()),
        StructField('sourceids', ArrayType(IntegerType()))]))),
        StructField('first_name', StringType()),
        StructField('kare_id', StringType()),
        StructField('last_name', StringType()),
        StructField('match_key', ArrayType(StringType()))
    ])
    

    来自 PySpark Docs on Row:“行可用于通过使用命名参数创建行对象,字段将按名称排序。”

    https://spark.apache.org/docs/latest/api/python/pyspark.sql.html#pyspark.sql.Row

    【讨论】:

    • 非常感谢。干净整洁。非常感谢
    【解决方案2】:

    首先,您实际上是在创建数据时定义了两次架构,当时您已经在 RDD 中使用行对象,因此您不需要使用 createDataFrame 函数,而是可以执行以下操作:

    sc.parallelize(row).toDF().show()
    

    但是,如果您仍然想明确提及架构,那么您需要保持架构和数据的顺序相同,并且您提到的架构根据您传递的数据是不正确的。正确的架构是:

    schema = StructType([
      StructField('fields', ArrayType(StructType([StructField('source', StringType()),StructField('sourceids', ArrayType(IntegerType()))]))), 
      StructField('first_name',StringType()), 
      StructField('kare_id',StringType()),
      StructField('last_name',StringType()), 
      StructField('match_key',ArrayType(StringType()))
    ])
    

    kare_id 应该在 last_name 之前,因为这是您传递数据的顺序

    【讨论】:

      猜你喜欢
      • 2021-06-01
      • 2017-02-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-01-29
      • 1970-01-01
      • 2021-09-11
      • 1970-01-01
      相关资源
      最近更新 更多