【问题标题】:Implementing a recursive algorithm in pyspark to find pairings within a dataframe在 pyspark 中实现递归算法以在数据帧中查找配对
【发布时间】:2020-11-19 10:57:22
【问题描述】:

我有一个 spark 数据框 (prof_student_df),其中列出了学生/教授对的时间戳。每个时间戳有 4 位教授和 4 位学生,每个教授-学生对都有一个“分数”(因此每个时间框架有 16 行)。对于每个时间范围,我需要找到教授/学生之间的一对一配对,以最大限度地提高总分。每个教授在一个时间范围内只能与一个学生匹配。

例如,这是一个时间范围内的配对/得分。

+------------+--------------+------------+-------+----------+
|    time    | professor_id | student_id | score | is_match |
+------------+--------------+------------+-------+----------+
| 1596048041 | p1           | s1         |   0.7 | FALSE    |
| 1596048041 | p1           | s2         |   0.5 | TRUE     |
| 1596048041 | p1           | s3         |   0.3 | FALSE    |
| 1596048041 | p1           | s4         |   0.2 | FALSE    |
| 1596048041 | p2           | s1         |   0.9 | TRUE     |
| 1596048041 | p2           | s2         |   0.1 | FALSE    |
| 1596048041 | p2           | s3         |  0.15 | FALSE    |
| 1596048041 | p2           | s4         |   0.2 | FALSE    |
| 1596048041 | p3           | s1         |   0.2 | FALSE    |
| 1596048041 | p3           | s2         |   0.3 | FALSE    |
| 1596048041 | p3           | s3         |   0.4 | FALSE    |
| 1596048041 | p3           | s4         |   0.8 | TRUE     |
| 1596048041 | p4           | s1         |   0.2 | FALSE    |
| 1596048041 | p4           | s2         |   0.3 | FALSE    |
| 1596048041 | p4           | s3         |  0.35 | TRUE     |
| 1596048041 | p4           | s4         |   0.4 | FALSE    |
+------------+--------------+------------+-------+----------+

目标是得到这个 is_match 列。它可以是布尔值或 0/1 位或任何有效的值。

在上面的示例中,p1 与 s2 匹配,p2 与 s1 匹配,p3 与 s4 匹配,p4 与 s3 匹配,因为这是使总分最大化的组合(得分为 2.55)。 有一个奇怪的极端情况——在给定的时间范围内,教授或学生可能少于 4 人。如果有 4 位教授和 3 位学生,那么 1 位教授将没有配对,并且他的所有 is_match 都是错误的。同样,如果有 3 位教授和 4 位学生,则 1 位学生将没有配对,并且他的所有 is_match 都将为 false。

有谁知道我可以如何做到这一点?我在想我会按时间进行分区或分组,然后将数据输入到一些 UDF 中,该 UDF 会吐出配对,然后也许我必须将其加入到原始行中(尽管我不确定)。我正在尝试在 pyspark 中实现这个逻辑,并且可以使用 spark sql/sql 或 pyspark。

理想情况下,我希望它尽可能高效,因为会有数百万行。在问题中,我提到了递归算法,因为这是一个传统的递归类型问题,但如果有更快的解决方案不使用递归,我愿意接受。

非常感谢,我是 spark 新手,对如何做到这一点有点困惑。

编辑:澄清问题,因为我在我的示例中意识到我没有指定这一点 一天,将有多达 14 位教授和 14 位学生可供选择。我一次只看一天,这就是为什么我在数据框中没有日期。在任何一个时间范围内,最多有 4 位教授和 4 位学生。此数据框仅显示一个时间范围。但在下一个时间范围内,这 4 位教授可能是 p5p1p7p9 或类似的东西。学生可能仍然是s1s2s3s4

【问题讨论】:

  • 我只看到了两种解决方法,1)窗口函数与数组/高阶函数(spark2.4+)的组合。 2) 熊猫 udaf (spark2.3+)。您的逻辑需要时间范围内的行之间的通信(以确保最大分数结果并仅在一个时间范围内使用不同的学生 ID),并且任何一种方式都将是计算密集型的。我认为使用数组/高阶函数会变得太复杂,并且使用 pandas 分组地图 udaf 可能会更好。我的 2 美分
  • 不同组合的数量是否固定为16?
  • @murtihash 您对如何使用 pandas 分组地图 udaf 执行此操作有什么建议吗?
  • @cronoik - 每行最多有 4 名学生和 4 名教授,对于每一行,我们计算一个教授学生对的值。如果缺少教授/学生,可能会少于 16 个组合,但永远不会更多。
  • 你需要实现类似hungarian algorithm的东西。

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


【解决方案1】:

正如我们的朋友@cronoik 提到的,您需要使用匈牙利算法,我在python 中看到的不平衡分配 问题的最佳代码是: https://github.com/mayorx/hungarian-algorithm(在存储库中也有一些示例:))

您只需要将 DataFrame 转换为 Numpy 数组并传递给 KM_Matcher,然后在 spark 中添加一个具有 withColumn 函数的列,具体取决于您对 KM_Matcher 的回答。

【讨论】:

  • 如何将数据帧转换为 numpy 数组?这是使用 pyspark 执行此操作的最有效方法吗
【解决方案2】:

编辑: 正如 cmets 中所讨论的,要解决您更新中提到的问题,我们可以将 student_id 每次转换为使用 dense_rank 的广义序列 ID,执行步骤 1 到 3(使用 student列),然后在每个时间使用join将student转换回原来的student_id。见下文Step-0Step-4。如果一个 timeUnit 中的教授少于 4 个,Numpy-end 中的维度将调整为 4(使用 np_vstack() 和 np_zeros()),请参阅更新的函数 find_assigned

你可以试试pandas_udfscipy.optimize.linear_sum_assignment(注意:后端方法是主cmets中@cronoik提到的匈牙利算法),见下文:

from pyspark.sql.functions import pandas_udf, PandasUDFType, first, expr, dense_rank
from pyspark.sql.types import StructType
from scipy.optimize import linear_sum_assignment
from pyspark.sql import Window
import numpy as np

df = spark.createDataFrame([
    ('1596048041', 'p1', 's1', 0.7), ('1596048041', 'p1', 's2', 0.5), ('1596048041', 'p1', 's3', 0.3),
    ('1596048041', 'p1', 's4', 0.2), ('1596048041', 'p2', 's1', 0.9), ('1596048041', 'p2', 's2', 0.1),
    ('1596048041', 'p2', 's3', 0.15), ('1596048041', 'p2', 's4', 0.2), ('1596048041', 'p3', 's1', 0.2),
    ('1596048041', 'p3', 's2', 0.3), ('1596048041', 'p3', 's3', 0.4), ('1596048041', 'p3', 's4', 0.8),
    ('1596048041', 'p4', 's1', 0.2), ('1596048041', 'p4', 's2', 0.3), ('1596048041', 'p4', 's3', 0.35),
    ('1596048041', 'p4', 's4', 0.4)
] , ['time', 'professor_id', 'student_id', 'score'])

N = 4
cols_student = [*range(1,N+1)]

步骤 0: 添加一个额外的列 student,并创建一个新的数据框 df3,其中包含 time + student_id + student 的所有唯一组合。

w1 = Window.partitionBy('time').orderBy('student_id')

df = df.withColumn('student', dense_rank().over(w1))
+----------+------------+----------+-----+-------+                              
|      time|professor_id|student_id|score|student|
+----------+------------+----------+-----+-------+
|1596048041|          p1|        s1|  0.7|      1|
|1596048041|          p2|        s1|  0.9|      1|
|1596048041|          p3|        s1|  0.2|      1|
|1596048041|          p4|        s1|  0.2|      1|
|1596048041|          p1|        s2|  0.5|      2|
|1596048041|          p2|        s2|  0.1|      2|
|1596048041|          p3|        s2|  0.3|      2|
|1596048041|          p4|        s2|  0.3|      2|
|1596048041|          p1|        s3|  0.3|      3|
|1596048041|          p2|        s3| 0.15|      3|
|1596048041|          p3|        s3|  0.4|      3|
|1596048041|          p4|        s3| 0.35|      3|
|1596048041|          p1|        s4|  0.2|      4|
|1596048041|          p2|        s4|  0.2|      4|
|1596048041|          p3|        s4|  0.8|      4|
|1596048041|          p4|        s4|  0.4|      4|
+----------+------------+----------+-----+-------+

df3 = df.select('time','student_id','student').dropDuplicates()
+----------+----------+-------+                                                 
|      time|student_id|student|
+----------+----------+-------+
|1596048041|        s1|      1|
|1596048041|        s2|      2|
|1596048041|        s3|      3|
|1596048041|        s4|      4|
+----------+----------+-------+

第 1 步: 使用 pivot 来找到教授与学生的矩阵,注意我们将分数设置为 pivot 的值,以便我们可以使用 scipy.optimize.linear_sum_assignment 来找到最小值分配问题的成本:

df1 = df.groupby('time','professor_id').pivot('student', cols_student).agg(-first('score'))
+----------+------------+----+----+-----+----+
|      time|professor_id|   1|   2|    3|   4|
+----------+------------+----+----+-----+----+
|1596048041|          p4|-0.2|-0.3|-0.35|-0.4|
|1596048041|          p2|-0.9|-0.1|-0.15|-0.2|
|1596048041|          p1|-0.7|-0.5| -0.3|-0.2|
|1596048041|          p3|-0.2|-0.3| -0.4|-0.8|
+----------+------------+----+----+-----+----+

Step-2:使用pandas_udf和scipy.optimize.linear_sum_assignment获取列索引,然后将对应的列名分配给新列assigned

# returnSchema contains one more StringType column `assigned` than schema from the input pdf:
schema = StructType.fromJson(df1.schema.jsonValue()).add('assigned', 'string')

# since the # of students are always N, we can use np.vstack to set the N*N matrix
# below `n` is the number of professors/rows in pdf
# sz is the size of input Matrix, sz=4 in this example
def __find_assigned(pdf, sz):
  cols = pdf.columns[2:]
  n = pdf.shape[0]
  n1 = pdf.iloc[:,2:].fillna(0).values
  _, idx = linear_sum_assignment(np.vstack((n1,np.zeros((sz-n,sz)))))
  return pdf.assign(assigned=[cols[i] for i in idx][:n])

find_assigned = pandas_udf(lambda x: __find_assigned(x,N), schema, PandasUDFType.GROUPED_MAP)

df2 = df1.groupby('time').apply(find_assigned)
+----------+------------+----+----+-----+----+--------+
|      time|professor_id|   1|   2|    3|   4|assigned|
+----------+------------+----+----+-----+----+--------+
|1596048041|          p4|-0.2|-0.3|-0.35|-0.4|       3|
|1596048041|          p2|-0.9|-0.1|-0.15|-0.2|       1|
|1596048041|          p1|-0.7|-0.5| -0.3|-0.2|       2|
|1596048041|          p3|-0.2|-0.3| -0.4|-0.8|       4|
+----------+------------+----+----+-----+----+--------+

注意:根据@OluwafemiSule 的建议,我们可以使用参数maximize 而不是否定得分值。该参数可用SciPy 1.4.0+

  _, idx = linear_sum_assignment(np.vstack((n1,np.zeros((N-n,N)))), maximize=True)

Step-3: 使用 SparkSQL stack 函数对上述 df2 进行归一化,取反分值并过滤分值为 NULL 的行。所需的is_match 列应该有assigned==student:

df_new = df2.selectExpr(
  'time',
  'professor_id',
  'assigned',
  'stack({},{}) as (student, score)'.format(len(cols_student), ','.join("int('{0}'), -`{0}`".format(c) for c in cols_student))
) \
.filter("score is not NULL") \
.withColumn('is_match', expr("assigned=student"))

df_new.show()
+----------+------------+--------+-------+-----+--------+
|      time|professor_id|assigned|student|score|is_match|
+----------+------------+--------+-------+-----+--------+
|1596048041|          p4|       3|      1|  0.2|   false|
|1596048041|          p4|       3|      2|  0.3|   false|
|1596048041|          p4|       3|      3| 0.35|    true|
|1596048041|          p4|       3|      4|  0.4|   false|
|1596048041|          p2|       1|      1|  0.9|    true|
|1596048041|          p2|       1|      2|  0.1|   false|
|1596048041|          p2|       1|      3| 0.15|   false|
|1596048041|          p2|       1|      4|  0.2|   false|
|1596048041|          p1|       2|      1|  0.7|   false|
|1596048041|          p1|       2|      2|  0.5|    true|
|1596048041|          p1|       2|      3|  0.3|   false|
|1596048041|          p1|       2|      4|  0.2|   false|
|1596048041|          p3|       4|      1|  0.2|   false|
|1596048041|          p3|       4|      2|  0.3|   false|
|1596048041|          p3|       4|      3|  0.4|   false|
|1596048041|          p3|       4|      4|  0.8|    true|
+----------+------------+--------+-------+-----+--------+

步骤 4: 使用 join 将 student 转换回 student_id(如果可能,使用广播 join):

df_new = df_new.join(df3, on=["time", "student"])
+----------+-------+------------+--------+-----+--------+----------+            
|      time|student|professor_id|assigned|score|is_match|student_id|
+----------+-------+------------+--------+-----+--------+----------+
|1596048041|      1|          p1|       2|  0.7|   false|        s1|
|1596048041|      2|          p1|       2|  0.5|    true|        s2|
|1596048041|      3|          p1|       2|  0.3|   false|        s3|
|1596048041|      4|          p1|       2|  0.2|   false|        s4|
|1596048041|      1|          p2|       1|  0.9|    true|        s1|
|1596048041|      2|          p2|       1|  0.1|   false|        s2|
|1596048041|      3|          p2|       1| 0.15|   false|        s3|
|1596048041|      4|          p2|       1|  0.2|   false|        s4|
|1596048041|      1|          p3|       4|  0.2|   false|        s1|
|1596048041|      2|          p3|       4|  0.3|   false|        s2|
|1596048041|      3|          p3|       4|  0.4|   false|        s3|
|1596048041|      4|          p3|       4|  0.8|    true|        s4|
|1596048041|      1|          p4|       3|  0.2|   false|        s1|
|1596048041|      2|          p4|       3|  0.3|   false|        s2|
|1596048041|      3|          p4|       3| 0.35|    true|        s3|
|1596048041|      4|          p4|       3|  0.4|   false|        s4|
+----------+-------+------------+--------+-----+--------+----------+

df_new = df_new.drop("student", "assigned")

【讨论】:

  • 您可以将maximize 选项传递给linear_sum_assignment 并且没有明确使用负权重
  • 谢谢@OluwafemiSule,我在您的建议中添加了注释。我的服务器有 SciPy 版本 1.2.0,它不支持这个参数,所以只保留旧的逻辑。
  • @jxc 非常感谢您在这里的帮助,这太棒了,我感谢您的彻底回应,因为它帮助我完成了它。一个快速的问题,这可能是我没有澄清的错 - 我只是在问题中澄清了,如果有 4 位教授和 4 位学生并不总是相同,这个解决方案是否有效?例如,在连续的许多时间范围内,可能是相同的 4 位教授和 4 位学生,但随后可能是新教授 (p5) 或混合中的新学生。永远只有 4 位教授和 4 位学生,但最多有 14 位教授和 14 位学生
  • @jxc 我意识到我认为我没有澄清这一点/想知道它是否仍然有效的原因是因为我在步骤 1 中看到最后一部分我们得到了所有学生的列表,但是该列表将包含在特定时间范围内未考虑的学生
  • @LaurenLeder,我调整了 pandas_udf 函数以处理处理器数小于 4 时的问题。还有 NULL 值问题,从 4*4 矩阵馈送到 linear_sum_assignment 的所有缺失值都将为零.让我知道这是否适合您的任务。
猜你喜欢
  • 1970-01-01
  • 2022-01-08
  • 2012-01-31
  • 2021-09-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2010-12-04
  • 1970-01-01
相关资源
最近更新 更多