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

PySpark DataFrame左连接结果异常,未返回预期0值求助

问题分析与解决方案

你对左连接的结果预期存在偏差,当前连接条件下df_2的两行都能在df_1中找到匹配项,因此不会生成amount=0的未匹配行;结果中的3条记录是因为df_2的第二行在df_1中有2个匹配项,左连接后会将匹配到的多行全部展开。

匹配逻辑说明

  • df_2第一行(col3=3051):与df_1第三行的col1/col2/col3/period字段完全匹配,因此生成1条匹配记录。
  • df_2第二行(col3=0000):与df_1前两行的上述四个字段匹配,因此生成2条匹配记录。
  • 左连接会保留左表(df_2)的所有行,并关联右表(df_1)中所有符合条件的行,最终结果行数为1+2=3,而非你预期的2行未匹配记录。

解决方案

根据你的真实需求,分两种场景处理:

场景1:保留df_2的行数,聚合匹配结果

如果你希望结果行数与df_2一致,对匹配到的多条记录进行聚合(比如求和amount、取第一个非空col7),可以在连接后分组聚合:

from pyspark.sql import functions as F

df = df_2.join(df_1
               , on = ["col1", "col2", "col3", "period"]
               , how = 'left') \
         .groupBy(df_2["col1"], df_2["col2"], df_2["col3"]) \
         .agg(
             F.first(df_1["col7"], ignorenulls=True).alias("col7"),
             F.coalesce(F.sum(df_1["amount"]), F.lit(0)).alias("amount")
         )
df.show(truncate=False)

执行后结果:

+-----+----+----+----+----------+
|col1 |col2|col3|col7|amount    |
+-----+----+----+----+----------+
|20110|b925|3051|7588|1256508.18|
|20110|b925|0000|8200|7         |
+-----+----+----+----+----------+

场景2:补充连接条件,生成未匹配行

如果你确实需要得到amount=0的未匹配行,需要补充更多连接字段(比如col4、col5),让df_2的行无法与df_1匹配:

from pyspark.sql import functions as F

df = df_2.join(df_1
               , on = (
                            (df_2["col3"] == df_1["col3"])
                          & (df_2["col1"] == df_1["col1"])
                          & (df_2["col2"] == df_1["col2"])
                          & (df_2["period"] == df_1["period"])
                          & (df_2["col4"] == df_1["col4"])
                          & (df_2["col5"] == df_1["col5"])
                        )
               , how = 'left').select(
                                df_2["col1"],
                                df_2["col2"],
                                df_2["col3"],
                                df_1["col7"],
                                F.coalesce(df_1["amount"], F.lit(0)).alias("amount")
                                )
df.show(truncate=False)

执行后结果:

+-----+----+----+----+------+
|col1 |col2|col3|col7|amount|
+-----+----+----+----+------+
|20110|b925|3051|null|0.0   |
|20110|b925|0000|null|0.0   |
+-----+----+----+----+------+

内容的提问来源于stack exchange,提问作者djm

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 03:22:49