You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于第三个PySpark DataFrame关联另外两个PySpark DataFrame

需求:通过Control DataFrame选择对应数据源的数据

我希望利用第三个PySpark DataFrame(control)来关联first和second两个PySpark DataFrame,control中记录了每条数据应取自前两个中的哪一个。


原始DataFrames

first DataFrame

ID1ID2ID3Data1Data2
203412444100200
20341223316332546400
321111311456544
311313441333645

second DataFrame

ID1ID2ID3Data1Data2
203412444133444
20341223333334211
3211113117685867443
311313441654463457

control DataFrame

ID2SOURCE
11first
12second
13first

期望结果

ID1ID2ID3Data1Data2
203412444133444
20341223333334211
321111311456544
311313441333645

现有代码框架

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 00:01:06