【问题标题】:PySpark - Join two Data Frames on Array column (order does not matter)PySpark - 在数组列上加入两个数据框(顺序无关紧要)
【发布时间】:2019-11-05 05:36:30
【问题描述】:

我在将两个 Dataframe 与 PySpark 中包含数组的列连接时遇到问题。如果数组中的元素相同(顺序无关紧要),我想加入这些列。

所以,我有一个 DataFrame,其中包含以下格式的项集及其频率:

+--------------------+----+
|               items|freq|
+--------------------+----+
|  [1828545, 1242385]|   4|
|  [1828545, 2032007]|   4|
|           [1137808]|  11|
|           [1209448]|   5|
|             [21002]|   5|
|           [2793224]| 209|
|     [2793224, 8590]|   7|
|[2793224, 8590, 8...|   4|
|[2793224, 8590, 8...|   4|
|[2793224, 8590, 8...|   5|
|[2793224, 8590, 1...|   4|
|  [2793224, 2593971]|  20|
+--------------------+----+

还有另一个 DataFrame,其中包含有关用户和项目的信息,格式如下:

+------------+-------------+--------------------+
|     user_id|   session_id| itemset            |
+------------+-------------+--------------------+
|WLB2T1JWGTHH|0012c5936056e|[1828545, 1242385]  |
|BZTAWYQ70C7N|00783934ea027|[2793224, 8590]     | 
|42L1RJL436ST|00c6821ed171e|[8590, 2793224]     |
|HB348HWSJAOP|00fa9607ead50|[21002]             |
|I9FOENUQL1F1|013f69b45bb58|[21002]             |  
+------------+-------------+--------------------+

现在,如果数组中的元素相同,我想在 itemset 和 items 上加入这两个数据框(它们的排序方式无关紧要)。我想要的输出是:

+------------+-------------+--------------------+----+
|     user_id|   session_id| itemset            |freq|
+------------+-------------+--------------------+----+
|WLB2T1JWGTHH|0012c5936056e|[1828545, 1242385]  |   4|
|BZTAWYQ70C7N|00783934ea027|[2793224, 8590]     |   7|
|42L1RJL436ST|00c6821ed171e|[8590, 2793224]     |   7|
|HB348HWSJAOP|00fa9607ead50|[21002]             |   5|
|I9FOENUQL1F1|013f69b45bb58|[21002]            |   5|  
+------------+-------------+--------------------+----+

我在网上找不到任何解决方案,只有在数组中包含一个项目时加入数据框的解决方案。

非常感谢! :)

【问题讨论】:

    标签: apache-spark dataframe join pyspark rdd


    【解决方案1】:

    join 的 spark 实现可以毫无问题地处理数组列。唯一的问题是,它不会忽略列的顺序。因此需要对连接列进行排序才能正确连接。您可以为此使用sort_array 函数。

    from pyspark.sql import functions as F
    
    df1 = spark.createDataFrame(
    [
    (  [1828545, 1242385],   4),
    (  [1828545, 2032007],   4),
    (           [1137808],  11),
    (           [1209448],   5),
    (             [21002],   5),
    (           [2793224], 209),
    (     [2793224, 8590],   7),
    ([2793224, 8590, 81],   4),
    ([2793224, 8590, 82],   4),
    ([2793224, 8590, 83],   5),
    ([2793224, 8590, 11],   4),
    (  [2793224, 2593971],  20)
    ], ['items','freq'])
    
    
    df2 = spark.createDataFrame(
    [
    ('WLB2T1JWGTHH','0012c5936056e',[1828545, 1242385]  ),
    ('BZTAWYQ70C7N','00783934ea027',[2793224, 8590]     ), 
    ('42L1RJL436ST','00c6821ed171e',[8590, 2793224]     ),
    ('HB348HWSJAOP','00fa9607ead50',[21002]             ),
    ('I9FOENUQL1F1','013f69b45bb58',[21002]             ) 
    ], ['user_id',   'session_id', 'itemset'])
    
    df1 = df1.withColumn('items', F.sort_array('items'))
    df2 = df2.withColumnRenamed('itemset', 'items').withColumn('items', F.sort_array('items'))
    
    df1.join(df2, "items").show()
    

    输出:

    +------------------+----+------------+-------------+ 
    |             items|freq|     user_id|   session_id| 
    +------------------+----+------------+-------------+ 
    |   [8590, 2793224]|   7|BZTAWYQ70C7N|00783934ea027| 
    |   [8590, 2793224]|   7|42L1RJL436ST|00c6821ed171e| 
    |[1242385, 1828545]|   4|WLB2T1JWGTHH|0012c5936056e| 
    |           [21002]|   5|HB348HWSJAOP|00fa9607ead50| 
    |           [21002]|   5|I9FOENUQL1F1|013f69b45bb58| 
    +------------------+----+------------+-------------+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-11-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-05-03
      • 2016-06-17
      相关资源
      最近更新 更多