如何基于第三个PySpark DataFrame关联另外两个PySpark DataFrame
需求:通过Control DataFrame选择对应数据源的数据
我希望利用第三个PySpark DataFrame(control)来关联first和second两个PySpark DataFrame,control中记录了每条数据应取自前两个中的哪一个。
原始DataFrames
first DataFrame
| ID1 | ID2 | ID3 | Data1 | Data2 |
|---|---|---|---|---|
| 2034 | 12 | 444 | 100 | 200 |
| 2034 | 12 | 233 | 1633 | 2546400 |
| 3211 | 11 | 311 | 456 | 544 |
| 3113 | 13 | 441 | 333 | 645 |
second DataFrame
| ID1 | ID2 | ID3 | Data1 | Data2 |
|---|---|---|---|---|
| 2034 | 12 | 444 | 133 | 444 |
| 2034 | 12 | 233 | 333 | 34211 |
| 3211 | 11 | 311 | 7685 | 867443 |
| 3113 | 13 | 441 | 6544 | 63457 |
control DataFrame
| ID2 | SOURCE |
|---|---|
| 11 | first |
| 12 | second |
| 13 | first |
期望结果
| ID1 | ID2 | ID3 | Data1 | Data2 |
|---|---|---|---|---|
| 2034 | 12 | 444 | 133 | 444 |
| 2034 | 12 | 233 | 333 | 34211 |
| 3211 | 11 | 311 | 456 | 544 |
| 3113 | 13 | 441 | 333 | 645 |
现有代码框架
cols_list = [ # cols with aliases to choose ] first = first.alias("a").join( second.alias("b"), ((first['ID1'] == second['ID1']) & first['ID2'] == second['ID2']) & first['ID3'] == second['ID3'])), 'left' ).select(cols_list)
实现方案
方法一:筛选后合并(高效优先)
这种方式先根据control的规则分别筛选出first和second中符合条件的数据,再合并结果,避免不必要的全量关联,性能更优。
from pyspark.sql import functions as F # 筛选first中对应SOURCE为'first'的行 first_selected = first.join(control, on="ID2", how="inner") \ .filter(F.col("SOURCE") == "first") \ .drop("SOURCE") # 筛选second中对应SOURCE为'second'的行 second_selected = second.join(control, on="ID2", how="inner") \ .filter(F.col("SOURCE") == "second") \ .drop("SOURCE") # 合并两个筛选后的数据集 final_df = first_selected.unionByName(second_selected)
方法二:全关联后条件选列(灵活优先)
如果需要保留first和second的所有关联行,或者需要处理部分缺失场景,可以先全量关联两个数据源,再根据control的规则选择对应列:
from pyspark.sql import functions as F # 关联first和second,以ID1/ID2/ID3为关联键 combined = first.alias("a").join(second.alias("b"), on=["ID1", "ID2", "ID3"], how="inner") # 关联control数据源 combined_with_control = combined.join(control, on="ID2", how="inner") # 根据SOURCE字段选择对应的Data1和Data2 final_df = combined_with_control.select( "ID1", "ID2", "ID3", F.when(F.col("SOURCE") == "first", F.col("a.Data1")).otherwise(F.col("b.Data1")).alias("Data1"), F.when(F.col("SOURCE") == "first", F.col("a.Data2")).otherwise(F.col("b.Data2")).alias("Data2") )
执行任意一种方法后,final_df就是你需要的结果。
内容的提问来源于stack exchange,提问作者dawid2312
相关产品推荐
相关产品推荐

