【问题标题】:Pyspark joining of two dataframes results with error of duplicated values两个数据帧的 Pyspark 连接导致重复值错误
【发布时间】:2021-06-30 16:55:41
【问题描述】:

在加入两个数据帧时,我在 pyspark 中遇到问题。第一个数据框是单列数据框“zipcd”,第二个数据框是四列数据框。

每当我尝试加入两个数据帧时都会出现问题,因为 Pyspark 在我的新数据帧中返回我,关于 zipcd 的单列,它的所有值都相同的列(第一行在所有行中重复,而且不是这样)。

例如:

Zip.select("Zip").show()
+------------+
|         Zip|
+------------+
| 6.0651002E8|
| 6.0623002E8|
| 6.0077203E8|
| 6.0626528E8|
| 6.0077338E8|
|         0.0|

另一个数据框是zipcd:

zip_cd1.show()
+-----+
|zipcd|
+-----+
|60651|
|60623|
|60077|
|60626|
|60077|
|    0|

每当我尝试加入数据框时,总是会发生以下情况:

Zip1=zip_cd1.join(Zip).select('Zip','zipcd')
Zip1.show()
+------------+-----+
|         Zip|zipcd|
+------------+-----+
| 6.0651002E8|60651|
| 6.0623002E8|60651|
| 6.0077203E8|60651|
| 6.0626528E8|60651|
| 6.0077338E8|60651|
|         0.0|60651|

无论我是否更改联接类型都会发生这种情况,而且我不知道发生了什么。

预期输出:

+------------+-----+
|         Zip|zipcd|
+------------+-----+
| 6.0651002E8|60651|
| 6.0623002E8|60623|
| 6.0077203E8|60077|
| 6.0626528E8|60626|
| 6.0077338E8|60077|
|         0.0|0    |

【问题讨论】:

  • 预期输出是什么?
  • 我会更新我的问题
  • 为你的加入提供密钥,否则无论你使用什么加入,它都将是笛卡尔加入。
  • 没有可用的键,因为其中一个数据框只有一列(单列数据框)

标签: join pyspark


【解决方案1】:

如果两个数据帧的分区数和行数相同,您可以使用RDD.zip,然后根据结果重新创建数据帧:

zipped_rdd = zip.rdd.zip(zipcd.rdd).map(lambda x: (x[0]['Zip'], x[1]['zipcd']))
df = spark.createDataFrame(zipped_rdd, schema=['Zip', 'zipcd'])

如果数据帧具有不同数量的分区或行,则可以使用RDD.zipWithIndexfull outer join,如果zip

zipped_rdd = zip.rdd.zipWithIndex().map(lambda x: (x[1], x[0])).fullOuterJoin(
  zipcd.rdd.zipWithIndex().map(lambda x: (x[1], x[0])) 
).map(lambda x: (x[1][0]['Zip'] if x[1][0] != None else None, x[1][1]['zipcd'] if x[1][1] != None else None))
df = spark.createDataFrame(zipped_rdd, schema=['Zip', 'zipcd'])

结果:

+-----------+-----+
|        Zip|zipcd|
+-----------+-----+
|6.0651002E8|60651|
|6.0623002E8|60623|
|6.0077203E8|60077|
|6.0626528E8|60626|
|6.0077338E8|60077|
|        0.0|    0|
+-----------+-----+

【讨论】:

  • 你的效果如何?我尝试了该代码,它返回我只能使用具有相同分区数的 rdd 进行压缩(第一个 rdd zip 有 200 个分区,第二个只有 1 个)。我将第一个 rdd 重新分区为仅 1 个分区,然后它返回错误“只能压缩每个分区中具有相同数量元素的 RDD”您是如何避免这些错误的?
  • @ArnoldoOliva 测试数据在我看来就像两个数据框非常相似。我为更一般的情况添加了一个可能的解决方案。
【解决方案2】:

我以邮政编码 60651 存储为 6.0651002E8 的方式理解数据,基于此假设,我给出了以下解决方案。

    >>> df = spark.read.text('test.dat')
    >>> df.show()
    +-----------+
    |      value|
    +-----------+
    |6.0651002E8|
    |6.0623002E8|
    |6.0077203E8|
    |6.0626528E8|
    |6.0077338E8|
    |        0.0|
    +-----------+
    df.createOrReplaceTempView("Tempview")
    df = spark.sql("select substring(cast(cast(value as decimal) as string),1,5) zipcd,value as zip from TempView")
    +-----+-----------+
    |zipcd|        zip|
    +-----+-----------+
    |60651|6.0651002E8|
    |60623|6.0623002E8|
    |60077|6.0077203E8|
    |60626|6.0626528E8|
    |60077|6.0077338E8|
    |    0|        0.0|
    +-----+-----------+
    import pyspark.sql.functions as f
    from pyspark.sql.functions import *
    from pyspark.sql.types import *
    
    data = [
    ("60651","State1","County1"),
    ("60623","State2","County2"),
    ("60077","State2","County3"),
    ("60626","State2","County4"),
    ("60077","State2","County5"),
    ("0","Unkown","Unkown")
            ]
    
    schema = StructType([
    StructField('zip', StringType(),True),\
    StructField('state', StringType(),True),\
    StructField('county', StringType(),True),\
    ])
    df1 = spark.createDataFrame(data=data, schema=schema)
    df1.show()
    df.createOrReplaceTempView("Data")
    df1.createOrReplaceTempView("Data1")
>>> spark.sql("select * from Data1 A inner join Data B on A.zip=b.zipcd").show()
+-----+------+-------+-----+-----------+
|  zip| state| county|zipcd|        zip|
+-----+------+-------+-----+-----------+
|60651|State1|County1|60651|6.0651002E8|
|60623|State2|County2|60623|6.0623002E8|
|60077|State2|County3|60077|6.0077338E8|
|60077|State2|County3|60077|6.0077203E8|
|60626|State2|County4|60626|6.0626528E8|
|60077|State2|County5|60077|6.0077338E8|
|60077|State2|County5|60077|6.0077203E8|
|    0|Unkown| Unkown|    0|        0.0|
+-----+------+-------+-----+-----------+

【讨论】:

    猜你喜欢
    • 2017-11-02
    • 2016-09-16
    • 1970-01-01
    • 2016-10-30
    • 2020-02-13
    • 2020-03-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多