【问题标题】:Spark SQL group, map reduceSpark SQL 组,map reduce
【发布时间】:2021-05-21 10:45:42
【问题描述】:

我有以下名为“数据”的数据集:

+---------+-------------+------+
|   name  |      subject| mark |
+---------+-------------+------+
|     Anna|         math|    80|
|     Vlad|      history|    67|
|     Jack|          art|    78|
|    David|         math|    71|
|   Monica|          art|    65|
|     Alex|          lit|    59|
|     Mark|         math|    82|
+---------+-------------+------+

我想做一个 map-reduce 工作。

结果显示如下或类似:

Anna, David : 1
Anna, Mark : 1
David, mark: 1
Vlad, None : 1
Jack, Monica: 1
Alex, None : 1

我已尝试执行以下操作:

new_data = data.select(['name', 'subject']).show()

+---------+-------------+
|   name  |      subject| 
+---------+-------------+
|     Anna|         math|  
|     Vlad|      history|  
|     Jack|          art|   
|    David|         math|   
|   Monica|          art|    
|     Alex|          lit|    
|     Mark|         math|    
+---------+-------------+


data_new.groupBy('name','subject').count().show(10)

但是,这个命令并没有提供我需要的东西。

【问题讨论】:

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


    【解决方案1】:

    您可以使用主题进行自左连接,获取不同的对,然后添加 1 列。

    import pyspark.sql.functions as F
    
    result = df.alias('t1').join(df.alias('t2'),
        F.expr("t1.subject = t2.subject and t1.name != t2.name"), 
        'left'
    ).select(
        F.concat_ws(
            ', ',
            F.greatest('t1.name', F.coalesce('t2.name', F.lit('None'))),
            F.least('t1.name', F.coalesce('t2.name', F.lit('None')))
        ).alias('pair')
    ).distinct().withColumn('val', F.lit(1))
    
    result.show()
    +------------+---+
    |        pair|val|
    +------------+---+
    |  Alex, None|  1|
    | Anna, David|  1|
    |  Anna, Mark|  1|
    |  None, Vlad|  1|
    | David, Mark|  1|
    |Jack, Monica|  1|
    +------------+---+
    

    【讨论】:

    • @user14913431 尝试编辑后的答案。不幸的是,array_sort 仅适用于 spark > 2.4。
    【解决方案2】:

    过程可能是:

    1. 将具有相同学科的学生分组到一个数组中
    2. 调用udf 函数来创建数组项排列
    3. 添加一列显示每个主题的数字
    4. 调用explode函数为数组中的每一项创建单独的3行

    让我们一步一步来做: 第 1 步:分组

    import pyspark.sql.functions as F
    
    grouped_df = data_new.groupBy('subject').agg(F.collect_set('name').alias('students_array'))
    

    第二步:udf函数

    from itertools import permutations
    def permutatoin(df_col):
        result = sorted([e for e in set(permutations(df_col))])
        return result 
    spark.udf.register("perWithPython", permutatoin)
    grouped_df  = grouped_df.select('*', permutatoin('students_array'))
    

    第 3 步:为每个主题创建一个新的数字值列

    grouped_df = grouped_df .withColumn('subject_no', F.rowNumber().over(Window.partitionBy('subject'))
    

    第 4 步:创建单独的行

    grouped_df.select(grouped_df.subject_no, explode(grouped_df.students_array)).show(truncate=False)
    

    【讨论】:

    • NameError: name 'perWithPython' is not defined
    • 哎呀,对不起,我调用了sql函数,我更新了代码来调用python函数
    猜你喜欢
    • 2015-08-13
    • 2019-06-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多