PySpark关联两个DataFrame,实现第二个DF全量行匹配第一个DF每个id
PySpark 全量匹配双DataFrame实现方案
核心逻辑
你需要的是第一个DataFrame每一行匹配第二个DataFrame全部行,本质是笛卡尔积关联,完成关联后按规则重命名字段、计算目标值即可。
完整可运行代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("cross_join_demo").getOrCreate() # 构造第一个DF数据 df1_data = [(1, "H234", 3), (2, "H123", 4)] df1 = spark.createDataFrame(df1_data, schema=["id", "user", "score"]) # 重命名score为condition df1 = df1.withColumnRenamed("score", "condition") # 构造第二个DF数据 df2_data = [(1, "Blood Pressure", 2), (2, "Stroke", 4), (3, "Joint Pain", 3)] df2 = spark.createDataFrame(df2_data, schema=["id", "trait", "conditional_score"]) # 笛卡尔积关联两个DF join_df = df1.crossJoin(df2) # 按规则计算新的conditional_score,同时选择指定输出字段 result_df = join_df.select( df1["id"], # 取第一个DF的id作为输出id,若需要取第二个DF的id可修改为df2["id"] "user", "condition", "trait", F.when( F.col("condition") <= F.col("conditional_score"), F.col("condition") + F.col("conditional_score") ).otherwise(F.col("condition")).alias("conditional_score") ) # 打印结果 result_df.show(truncate=False)
输出结果验证
运行后输出如下,完全符合预期规则:
+---+----+---------+--------------+-------------------+ |id |user|condition|trait |conditional_score | +---+----+---------+--------------+-------------------+ |1 |H234|3 |Blood Pressure|3 | |1 |H234|3 |Stroke |7 | |1 |H234|3 |Joint Pain |6 | |2 |H123|4 |Blood Pressure|4 | |2 |H123|4 |Stroke |8 | |2 |H123|4 |Joint Pain |4 | +---+----+---------+--------------+-------------------+
注意事项
笛卡尔积会使数据量变为df1行数 * df2行数,如果两个DF数据量极大,需要提前评估集群资源是否足够。
内容的提问来源于stack exchange,提问作者Evandro Lippert
相关产品推荐
相关产品推荐

