PySpark条件关联咨询:含可空列的两表动态连接优化方案
问题
现有一段PySpark代码,将spark_company_df与spark_geo_df通过三列左外连接:
final_df = spark_company_df .join( spark_geo_df, (spark_company_df.column1 == spark_geo_df.column1) & (spark_company_df.column2 == spark_geo_df.column2) & (spark_company_df.column3 == spark_geo_df.column3), "left_outer") .select( spark_geo_df.column1, spark_geo_df.column2, spark_geo_df.column3)
其中spark_company_df的column3可为空,其余两列非空。当前逻辑下,当column3为空时,结果三列全为null,但期望此时仅通过前两列关联,得到如下期望结果:
-----------------+--------------------+-------------+- |column1 | column2 | column3| +-----------------+--------------------+-------------+ | LA| CA| US| | LA| CA| US| | SF| CA| null| +-----------------+--------------------+-------------+
而实际得到的结果是:
-----------------+--------------------+-------------+- |column1 | column2 | column3| +-----------------+--------------------+-------------+ | LA| CA| US| | LA| CA| US| | null| null| null| +-----------------+--------------------+-------------+
需要找到无需额外关联的条件连接方案,实现column3非空时用三列关联,为空时用两列关联。
解决方案
可以通过调整连接条件,判断spark_company_df.column3是否为空,为空时只匹配前两列,非空时匹配三列。具体代码如下:
final_df = spark_company_df.join( spark_geo_df, # 核心条件:column3为空时仅匹配前两列;非空时三列全匹配 (spark_company_df.column1 == spark_geo_df.column1) & (spark_company_df.column2 == spark_geo_df.column2) & ( spark_company_df.column3.isNull() | (spark_company_df.column3 == spark_geo_df.column3) ), "left_outer" ).select( spark_geo_df.column1, spark_geo_df.column2, spark_geo_df.column3 )
逻辑说明
- 当
spark_company_df.column3为空时,spark_company_df.column3.isNull()为True,此时只需满足前两列匹配即可关联spark_geo_df - 当
spark_company_df.column3非空时,spark_company_df.column3.isNull()为False,此时需要同时满足第三列匹配,即三列全匹配才会关联
这样就能在一次左外连接中实现需求,无需额外关联操作。
内容的提问来源于stack exchange,提问作者Anton Kim
相关产品推荐
相关产品推荐

