PySpark中如何对Join的每一列参数使用Coalesce?
解决Spark Leftanti Join含Null值的简便方法
直接用列表推导式批量生成带coalesce的匹配条件,不用逐个列编写:
from pyspark.sql.functions import coalesce, col attrs = ['surname', 'name', 'patronymic', 'birth_date', 'doc_type', 'doc_series','doc_number'] # 生成每列的coalesce匹配条件,用空字符串替代null(可根据业务换特殊值) join_conditions = [coalesce(pre_df[c], "") == coalesce(res_df[c], "") for c in attrs] # 执行leftanti join new_df = pre_df.join(res_df, on=join_conditions, how='leftanti')
说明
Spark中null值的比较结果为null,会被视为不匹配,导致join逻辑不符合预期。通过coalesce把null替换成一个业务中不会冲突的确定值(比如空字符串、特殊标记"__NULL__"),就能让含null的行正常参与匹配逻辑。
如果你希望保留原DataFrame的所有列,不需要提前做select处理,直接用上面的条件数组作为join的on参数即可。
内容的提问来源于stack exchange,提问作者schoolboychik
相关产品推荐
相关产品推荐

