如何在PySpark或SQL中实现多对多变量的一对一整行匹配
PySpark 解决方案
用窗口函数给分组内的行加序号,再关联,逻辑和SQL一致。
假设你已经有两个DataFrame:df1(含date, amount, description)和df2(含date, amount, ID):
from pyspark.sql import Window from pyspark.sql.functions import row_number # 给df1的每个(date, amount)分组加序号 window1 = Window.partitionBy("date", "amount").orderBy("description") df1_with_rank = df1.withColumn("row_num", row_number().over(window1)) # 给df2的每个(date, amount)分组加序号 window2 = Window.partitionBy("date", "amount").orderBy("ID") df2_with_rank = df2.withColumn("row_num", row_number().over(window2)) # 关联得到一对一匹配结果 matched_df = df1_with_rank.join( df2_with_rank, on=["date", "amount", "row_num"], how="inner" ) # 计算总额 total = matched_df.agg({"amount": "sum"}).collect()[0][0] print(total) # 预期输出1097
注意事项
- 排序字段(比如
orderBy("description"))可以随便选,只要同一分组内的序号唯一就行,不影响最终求和结果。 - 如果两个分组的记录数不一样,
innerjoin会取两者的最小记录数;要是需要保留所有记录,可换成left或fulljoin,但按你的需求inner足够。
内容的提问来源于stack exchange,提问作者Heather
相关产品推荐
相关产品推荐

